*
* backup.cpp: Backup api used by Backup/Recovery manager.
*
* Portions Copyright (c) 2020 Huawei Technologies Co.,Ltd.
* Portions Copyright (c) 2009-2013, NIPPON TELEGRAPH AND TELEPHONE CORPORATION
* Portions Copyright (c) 2015-2018, Postgres Professional
*
*-------------------------------------------------------------------------
*/
#include "include/backup.h"
#include "storage/file/fio_device.h"
#include "common/fe_memutils.h"
static int g_doneFiles = 0;
static int g_totalFiles = 0;
static volatile bool g_progressFlag = false;
static pthread_cond_t g_cond = PTHREAD_COND_INITIALIZER;
static pthread_mutex_t g_mutex = PTHREAD_MUTEX_INITIALIZER;
static void handleZeroSizeFile(FileAppender *appender, pgFile* file);
static void *ProgressReportProbackup(void *arg);
void performBackup(backup_files_arg* arg)
{
backupReaderThreadArgs* thread_args = (backupReaderThreadArgs*)palloc(sizeof(backupReaderThreadArgs) * current.readerThreadCount);
initPerformBackup(arg, thread_args);
backupDataFiles(arg);
pfree(thread_args);
}
void initPerformBackup(backup_files_arg* arg, backupReaderThreadArgs* thread_args)
{
startBackupSender();
startBackupReaders(arg, thread_args);
}
void backupDataFiles(backup_files_arg* arg)
{
FileAppender* appender = NULL;
* instance_name complies with the naming rules of object storage service
*/
appender = (FileAppender*)palloc(sizeof(FileAppender));
if (appender == NULL) {
elog(ERROR, "Failed to allocate memory for appender");
return;
}
appender->baseFileName = pg_strdup(current.root_dir);
pthread_spin_init(&appender->lock, PTHREAD_PROCESS_PRIVATE);
initFileAppender(appender, FILE_APPEND_TYPE_FILES, 0, 0);
backupDirectories(appender, arg);
backupFiles(appender, arg);
flushReaderContexts(arg);
closeFileAppender(appender);
destoryFileAppender(&appender);
stopBackupReaders();
destoryBackupReaderContexts();
setSenderState(current.sender_cxt, SENDER_THREAD_STATE_FINISH);
waitForSenderThread();
stopBackupSender();
destoryBackupSenderContext();
}
void backupDirectories(FileAppender* appender, backup_files_arg* arg)
{
int totalFiles = (int)parray_num(arg->files_list);
g_totalFiles = totalFiles;
for (int i = 0; i < totalFiles; i++) {
pgFile* dir = (pgFile*) parray_get(arg->files_list, i);
if (S_ISDIR(dir->mode)) {
char dirpath[MAXPGPATH];
int nRet = snprintf_s(dirpath, MAXPGPATH, MAXPGPATH - 1, "%s", dir->rel_path);
securec_check_ss_c(nRet, "\0", "\0");
appendDir(appender, dirpath, DIR_PERMISSION, dir->external_dir_num, dir->type);
g_doneFiles++;
}
}
}
void appendDir(FileAppender* appender, const char* dirPath, uint32 permission,
int external_dir_num, device_type_t type)
{
ParallelFileAppenderSegHeader header;
size_t pathLen = strlen(dirPath);
setSegHeaderVersion(&header, SEG_HEADER_PARALLEL_VERSION);
SetSegHeaderType(&header, FILE_APPEND_TYPE_DIR);
header.size = pathLen;
header.permission = permission;
header.filesize = 0;
header.crc = 0;
header.external_dir_num = external_dir_num;
header.file_type = type;
header.threadId = 0;
WriteDataBlock(&header, (char*)dirPath, pathLen, appender);
}
void backupFiles(FileAppender* appender, backup_files_arg* arg)
{
char from_fullpath[MAXPGPATH];
char to_fullpath[MAXPGPATH];
static time_t prev_time;
time_t start_time, end_time;
char pretty_time[20];
int n_backup_files_list = parray_num(arg->files_list);
prev_time = current.start_time;
parray_qsort(arg->files_list, pgFileCompareSize);
if (arg->prev_filelist)
parray_qsort(arg->prev_filelist, pgFileCompareRelPathWithExternal);
write_backup_filelist(¤t, arg->files_list,
instance_config.pgdata, arg->external_dirs, true);
write_backup(¤t, true);
init_header_map(¤t);
thread_interrupted = false;
elog(INFO, "Start backing up files");
time(&start_time);
pthread_t progressThread;
pthread_create(&progressThread, nullptr, ProgressReportProbackup, nullptr);
for (int i = 0; i < n_backup_files_list; i++) {
pgFile *prev_file = NULL;
pgFile *file = (pgFile *) parray_get(arg->files_list, i);
if (S_ISDIR(file->mode)) {
continue;
}
if (interrupted || thread_interrupted) {
elog(ERROR, "interrupted during backup");
}
if (progress)
elog_file(INFO, "Progress: (%d/%d). Process file \"%s\"",
i + 1, n_backup_files_list, file->rel_path);
pg_atomic_add_fetch_u32((volatile uint32*) &g_doneFiles, 1);
if (file->size == 0) {
file->write_size = 0;
handleZeroSizeFile(appender, file);
continue;
}
if (file->external_dir_num != 0) {
char external_dst[MAXPGPATH];
char *external_path = (char *)parray_get(arg->external_dirs,
file->external_dir_num - 1);
makeExternalDirPathByNum(external_dst,
arg->external_prefix,
file->external_dir_num);
join_path_components(to_fullpath, external_dst, file->rel_path);
join_path_components(from_fullpath, external_path, file->rel_path);
} else if (is_dss_type(file->type)) {
join_path_components(from_fullpath, arg->src_dss, file->rel_path);
join_path_components(to_fullpath, arg->dst_dss, file->rel_path);
} else {
join_path_components(from_fullpath, arg->from_root, file->rel_path);
join_path_components(to_fullpath, arg->to_root, file->rel_path);
}
if (!S_ISREG(file->mode)) {
elog(WARNING, "Unexpected type %d of file \"%s\", skipping",
file->mode, from_fullpath);
}
if (current.backup_mode != BACKUP_MODE_FULL) {
pgFile **prev_file_tmp = NULL;
prev_file_tmp = (pgFile **) parray_bsearch(arg->prev_filelist,
file, pgFileCompareRelPathWithExternal);
if (prev_file_tmp) {
file->exists_in_prev = true;
prev_file = *prev_file_tmp;
}
}
if (file->external_dir_num == 0 && strcmp(file->name, PG_XLOG_CONTROL_FILE) == 0) {
char* filename = last_dir_separator(to_fullpath);
char* dirpath = strndup(to_fullpath, filename - to_fullpath + 1);
fio_mkdir(dirpath, DIR_PERMISSION, FIO_BACKUP_HOST);
pg_free(dirpath);
}
* Select a free reader thread to backup files, probackup thread only dispatch tasks.
*/
int threadSlot = getFreeReaderThread();
while (threadSlot == -1) {
flushReaderContexts(arg);
threadSlot = getFreeReaderThread();
}
ReaderCxt* readerCxt = ¤t.readerCxt[threadSlot];
int currentFillIdx = readerCxt->fileCount;
readerCxt->file[currentFillIdx] = file;
readerCxt->prefile[currentFillIdx] = prev_file;
readerCxt->fromPath[currentFillIdx] = pgut_strdup(from_fullpath);
readerCxt->toPath[currentFillIdx] = pgut_strdup(to_fullpath);
readerCxt->appender = appender;
readerCxt->segType[currentFillIdx] = FILE_APPEND_TYPE_FILE;
readerCxt->fileRemoved[currentFillIdx] = false;
readerCxt->fileCount++;
if (readerCxt->fileCount == READER_THREAD_FILE_COUNT || file->size >= FILE_BUFFER_SIZE) {
setReaderState(readerCxt, READER_THREAD_STATE_START);
}
if (file->write_size == FILE_NOT_FOUND) {
continue;
}
if (file->write_size == BYTES_INVALID) {
elog(VERBOSE, "Skipping the unchanged file: \"%s\"", from_fullpath);
continue;
}
}
g_progressFlag = true;
pthread_mutex_lock(&g_mutex);
pthread_cond_signal(&g_cond);
pthread_mutex_unlock(&g_mutex);
pthread_join(progressThread, nullptr);
fio_disconnect();
if (arg->conn_arg.conn) {
pgut_disconnect(arg->conn_arg.conn);
}
arg->ret = 0;
elog(INFO, "Finish backuping file");
time(&end_time);
pretty_time_interval(difftime(end_time, start_time),
pretty_time, lengthof(pretty_time));
elog(INFO, "Backup files are backuped to s3, time elapsed: %s", pretty_time);
}
static void handleZeroSizeFile(FileAppender *appender, pgFile* file)
{
size_t pathLen = strlen(file->rel_path);
ParallelFileAppenderSegHeader startHeader;
constructParallelHeader(&startHeader, FILE_APPEND_TYPE_FILE, pathLen, 0, file, 0);
WriteDataBlock(&startHeader, (char*)file->rel_path, pathLen, appender);
ParallelFileAppenderSegHeader endHeader;
constructParallelHeader(&endHeader, FILE_APPEND_TYPE_FILE_END, 0, 0, file, 0);
WriteHeader(&endHeader, appender);
}
static void *ProgressReportProbackup(void *arg)
{
if (g_totalFiles == 0) {
return nullptr;
}
char progressBar[53];
int percent;
do {
percent = (int)(g_doneFiles * 100 / g_totalFiles);
GenerateProgressBar(percent, progressBar);
fprintf(stdout, "Progress: %s %d%% (%d/%d, done_files/total_files). backup file \r",
progressBar, percent, g_doneFiles, g_totalFiles);
pthread_mutex_lock(&g_mutex);
timespec timeout;
timeval now;
gettimeofday(&now, nullptr);
timeout.tv_sec = now.tv_sec + 1;
timeout.tv_nsec = 0;
int ret = pthread_cond_timedwait(&g_cond, &g_mutex, &timeout);
pthread_mutex_unlock(&g_mutex);
if (ret == ETIMEDOUT) {
continue;
} else {
break;
}
} while ((g_doneFiles < g_totalFiles) && !g_progressFlag);
percent = 100;
GenerateProgressBar(percent, progressBar);
fprintf(stdout, "Progress: %s %d%% (%d/%d, done_files/total_files). backup file \n",
progressBar, percent, g_doneFiles, g_totalFiles);
return nullptr;
}