file.cpp

来自「最新版本!fastdb是高效的内存数据库系统」· C++ 代码 · 共 1,638 行 · 第 1/4 页

CPP
1,638
字号
                return -1;
            }
            j = i;
        }
        if (i >= (size_t)nPages || currUpdateCount[i] > updateCounters[i]) { 
            if (currUpdateCount[i] > maxUpdateCount) { 
                maxUpdateCount = currUpdateCount[i];
            }
            updateCounters[i] = currUpdateCount[i];
        } else { 
            j = i + 1;
        }
    }      
    if (i != j) { 
        db->con[nodeId].nRecoveredPages += (i - j);
        rr.op = ReplicationRequest::RR_RECOVER_PAGE;
        rr.nodeId = nodeId;
        rr.size = (i-j)*dbModMapBlockSize;
        rr.page.offs = j << dbModMapBlockBits;
        rr.page.updateCount = currUpdateCount[j];
        TRACE_MSG(("Send segment [%lx, %ld]\n", rr.page.offs, rr.size));
        if (!db->writeReq(nodeId, rr, mmapAddr + rr.page.offs, rr.size)) {
            return -1;
        }
    }   
    return maxUpdateCount;
}

void dbFile::completeRecovery(int nodeId)
{
    ReplicationRequest rr;
    dbTrace("Complete recovery of node %d: recovere %d pages\n", nodeId, db->con[nodeId].nRecoveredPages);
    rr.op = ReplicationRequest::RR_STATUS;
    rr.nodeId = nodeId;
    db->con[nodeId].status = rr.status = dbReplicatedDatabase::ST_STANDBY;
    for (int i = 0, n = db->nServers; i < n; i++) {                 
        if (db->con[i].status != dbReplicatedDatabase::ST_OFFLINE && i != db->id) {
            db->writeReq(i, rr); 
        }
    }
}

void dbFile::doRecovery(int nodeId, int* updateCounters, int nPages)
{
    int maxUpdateCount;
    memset(updateCounters+nPages, 0, (getMaxPages() - nPages)*sizeof(int));

    if (db->con[nodeId].reqSock == NULL) { 
        char buf[256];
        socket_t* s = socket_t::connect(db->serverURL[nodeId], 
                                        socket_t::sock_global_domain, 
                                        db->recoveryConnectionAttempts);
        if (!s->is_ok()) { 
            s->get_error_text(buf, sizeof buf);
            dbTrace("Failed to establish connection with node %d: %s\n",
                    nodeId, buf);
            delete s;
            return;
        } 
        ReplicationRequest rr;
        rr.op = ReplicationRequest::RR_GET_STATUS;
        rr.nodeId = db->id;
        if (!s->write(&rr, sizeof rr) || !s->read(&rr, sizeof rr)) { 
            s->get_error_text(buf, sizeof buf);
            dbTrace("Connection with node %d is broken: %s\n",
                    nodeId, buf);
            delete s;
            return;
        }
        if (rr.op != ReplicationRequest::RR_STATUS && rr.status != dbReplicatedDatabase::ST_STANDBY) { 
            dbTrace("Unexpected response from standby node %d: code %d status %d\n", 
                     nodeId, rr.op, rr.status);
            delete s;
            return;
        } else {
            db->addConnection(nodeId, s);
        }
    }

    for (int i = 0; i < db->maxAsyncRecoveryIterations; i++) { 
        { 
            dbCriticalSection cs(syncCS);
            maxUpdateCount = sendChanges(nodeId, updateCounters, nPages);
            if (maxUpdateCount < 0) { 
                delete[] updateCounters;
                return;
            }
        }
        { 
            dbCriticalSection cs(replCS);
            if (maxUpdateCount == updateCounter) { 
                delete[] updateCounters;
                completeRecovery(nodeId);
                return;
            }
        }
    }
    {
        dbTrace("Syncronouse recovery of node %d\n", nodeId);
        dbCriticalSection cs1(syncCS);
        dbCriticalSection cs2(replCS); 
        maxUpdateCount = sendChanges(nodeId, updateCounters, nPages);
        delete[] updateCounters;
        if (maxUpdateCount >= 0) { 
            assert(maxUpdateCount == updateCounter);
            completeRecovery(nodeId);
        }
    }
}


#endif

bool dbFile::write(void const* buf, size_t size)
{
    size_t writtenBytes;
    bool result = write(buf, writtenBytes, size) == ok && writtenBytes == size;
    return result;
}

#ifdef _WIN32

class OS_info : public OSVERSIONINFO { 
  public: 
    OS_info() { 
        dwOSVersionInfoSize = sizeof(OSVERSIONINFO);
        GetVersionEx(this);
    }
};

static OS_info osinfo;

#define BAD_POS 0xFFFFFFFF // returned by SetFilePointer and GetFileSize


int dbFile::erase()
{
    return ok;
}


#ifdef PROTECT_DATABASE
void dbFile::protect(size_t pos, size_t size)
{
    PDWORD oldProtect;
    bool rc = VirtualProtect(mmapAddr + pos, DOALIGN(size, pageSize), PAGE_READONLY, &oldProtect);
    assert(rc);
}

void dbFile::unprotect(size_t pos, size_t size)
{
    PDWORD oldProtect;
    bool rc = VirtualProtect(mmapAddr + pos, DOALIGN(size, pageSize), PAGE_READWRITE, &oldProtect);
    assert(rc);    
}
#endif


int dbFile::open(char const* fileName, char const* sharedName, bool readonly,
                 size_t initSize, bool replicationSupport)
{
    int status;
    size_t fileSize = 0;
    this->readonly = readonly;
#ifndef DISKLESS_CONFIGURATION
    if (strcmp(fileName, "/dev/zero") == 0) { 
        fh = INVALID_HANDLE_VALUE;
        fileSize = initSize;
    } else { 
#if defined(_WINCE) && !defined(NO_MMAP)
        fh = CreateFileForMapping
#else
        fh = CreateFile
#endif
            (W32_STRING(fileName), 
             readonly ? GENERIC_READ : (GENERIC_READ|GENERIC_WRITE), 
             FILE_SHARE_READ | FILE_SHARE_WRITE, 
             FASTDB_SECURITY_ATTRIBUTES, 
             readonly ? OPEN_EXISTING : OPEN_ALWAYS,
#ifdef _WINCE
             FILE_ATTRIBUTE_NORMAL
#else
             FILE_FLAG_RANDOM_ACCESS
#ifdef NO_MMAP
             |FILE_FLAG_NO_BUFFERING
#endif
#if 0 // not needed as we do explicit flush ???
             |FILE_FLAG_WRITE_THROUGH
#endif
#endif
             , NULL);
        if (fh == INVALID_HANDLE_VALUE) {
            return GetLastError();
        }
        DWORD highSize;
        fileSize = GetFileSize(fh, &highSize);
        if (fileSize == BAD_POS && (status = GetLastError()) != ok) {
            CloseHandle(fh);
            return status;
        }
        assert(highSize == 0);
    }
    mmapSize = fileSize;

    this->sharedName = new char[strlen(sharedName) + 1];
    strcpy(this->sharedName, sharedName);

    if (!readonly && fileSize == 0) { 
        mmapSize = initSize;
    }
#else
    fh = INVALID_HANDLE_VALUE;
    this->sharedName = NULL;
    mmapSize = fileSize = initSize;
#endif
#if defined(NO_MMAP)
    if (fileSize < mmapSize && !readonly) { 
        if (SetFilePointer(fh, mmapSize, NULL, FILE_BEGIN) != mmapSize || !SetEndOfFile(fh)) {
            status = GetLastError();
            CloseHandle(fh);
            return status;
        }
    }
    mmapAddr = (char*)VirtualAlloc(NULL, mmapSize, MEM_COMMIT|MEM_RESERVE, 
                                   PAGE_READWRITE);
           
#ifdef DISKLESS_CONFIGURATION
    if (mmapAddr == NULL) 
#else
    size_t readBytes;
    if (mmapAddr == NULL
        || read(mmapAddr, readBytes, fileSize) != ok || readBytes != fileSize) 
#endif    
    {  
        status = GetLastError();
        if (fh != INVALID_HANDLE_VALUE) { 
            CloseHandle(fh);
        }
        return status;
    } 
    memset(mmapAddr+fileSize, 0, mmapSize - fileSize);
    mh = NULL;
#else
    mh = CreateFileMapping(fh, FASTDB_SECURITY_ATTRIBUTES, readonly ? PAGE_READONLY : PAGE_READWRITE, 
                           0, mmapSize, W32_STRING(sharedName));
    status = GetLastError();
    if (mh == NULL) { 
        if (fh != INVALID_HANDLE_VALUE) { 
            CloseHandle(fh);
        }
        return status;
    }
    mmapAddr = (char*)MapViewOfFile(mh, readonly ? FILE_MAP_READ : FILE_MAP_ALL_ACCESS, 0, 0, 0);
    if (mmapAddr == NULL) { 
        status = GetLastError();
        CloseHandle(mh);
        if (fh != INVALID_HANDLE_VALUE) { 
            CloseHandle(fh);
        }
        return status;
    } 
    if (status != ERROR_ALREADY_EXISTS && mmapSize > fileSize)
        //      && osinfo.dwPlatformId != VER_PLATFORM_WIN32_NT) 
    { 
        // Windows 95 doesn't initialize pages
        memset(mmapAddr+fileSize, 0, mmapSize - fileSize);
    }
#endif

#if defined(NO_MMAP) || defined(REPLICATION_SUPPORT)
    SYSTEM_INFO systemInfo;
    GetSystemInfo(&systemInfo);
    pageSize = systemInfo.dwPageSize;
    pageMapSize = (mmapSize + dbModMapBlockSize*32 - 1) >> (dbModMapBlockBits + 5);
    pageMap = new int[pageMapSize];
    memset(pageMap, 0, pageMapSize*sizeof(int));
#endif

#if defined(REPLICATION_SUPPORT)
    db = NULL;
    int nPages = getMaxPages();        
    currUpdateCount = new int[nPages];

    if (replicationSupport) { 
        char* cFileName = new char[strlen(fileName) + 5];
        strcat(strcpy(cFileName, fileName), ".cnt");
        
#ifdef DISKLESS_CONFIGURATION
        cfh = INVALID_HANDLE_VALUE;
#else
        cfh = CreateFile(cFileName, GENERIC_READ|GENERIC_WRITE, 
                         0, NULL, OPEN_ALWAYS,
                         FILE_FLAG_RANDOM_ACCESS|FILE_FLAG_WRITE_THROUGH,
                         NULL);
        delete[] cFileName;
        if (cfh == INVALID_HANDLE_VALUE) {
            status = errno;
            return status;
        }
#endif
        cmh = CreateFileMapping(cfh, NULL, PAGE_READWRITE, 0, 
                                nPages*sizeof(int), NULL);
        status = GetLastError();
        if (cmh == NULL) { 
            CloseHandle(cfh);
            return status;
        }
        diskUpdateCount = (int*)MapViewOfFile(cmh, FILE_MAP_ALL_ACCESS, 0, 0, 0);
        if (diskUpdateCount == NULL) { 
            status = GetLastError();
            CloseHandle(cmh);
            CloseHandle(cfh);
            return status;
        } 
        rootPage = dbMalloc(pageSize);
        int maxCount = 0;
        for (int i = 0; i < nPages; i++) {      
            int count = diskUpdateCount[i];
            currUpdateCount[i] = count;
            if (count > maxCount) { 
                maxCount = count;
            }
        }
        updateCounter = maxCount;
        nRecovered = 0;
        recoveredEvent.open(true);
        syncEvent.open();
        startSync();
    }
#endif

#ifdef FUZZY_CHECKPOINT
    writer = new dbFileWriter(this);
#endif

    return ok; 
}

bool dbFile::write(size_t pos, void const* ptr, size_t size) 
{
    DWORD written;
    if (SetFilePointer(fh, pos, NULL, FILE_BEGIN) != pos ||
        !WriteFile(fh, ptr, size, &written, NULL) 
        || written != (DWORD)size) 
    { 
        dbTrace("Failed to save page to the disk, position=%ld, size=%ld, error=%d\n",
                (long)pos, (long)size, GetLastError());
        return false;
    }
    return true;
}

#if defined(REPLICATION_SUPPORT)
void dbFile::syncToDisk()
{
    syncThread.setPriority(dbThread::THR_PRI_LOW);
    dbCriticalSection cs(syncCS);
    while (doSync) { 
        size_t i, j, k, n; 
        int maxUpdated = 0;
        for (i = 0, n = mmapSize >> dbModMapBlockBits; i < n;) { 
            int updateCounters[dbMaxSyncSegmentSize];
            for (j=i; j < (mmapSize >> dbModMapBlockBits) && j-i < dbMaxSyncSegmentSize 
                     && currUpdateCount[j] > diskUpdateCount[j]; j++)
            {
                updateCounters[j-i] = currUpdateCount[j];
            }
            if (i != j) { 
                size_t pos = (i << dbModMapBlockBits) & ~(pageSize-1);
                size_t size = (((j-i) << dbModMapBlockBits) + pageSize - 1) & ~(pageSize-1);
#ifdef NO_MMAP
                write(pos, mmapAddr + pos, size);
#else
                FlushViewOfFile(mmapAddr + pos, size);
#endif
                for (k = 0; i < j; k++, i++) {  
                    diskUpdateCount[i] = updateCounters[k];
                }
                maxUpdated = i;
            } else { 
                i += 1;
            }
            if (!doSync) { 
                return;
            }
        }
        if (maxUpdated != 0) { 
            FlushViewOfFile(diskUpdateCount, maxUpdated*sizeof(int));
        }
        if (closing && maxUpdated == 0) { 
            return;
        } else { 
            syncEvent.wait(syncCS, dbSyncTimeout);
        }
    }
}
#endif

int dbFile::create(const char* name, bool noBuffering)
{
    fh = CreateFile(W32_STRING(name), GENERIC_READ|GENERIC_WRITE, 0, FASTDB_SECURITY_ATTRIBUTES, CREATE_ALWAYS, 
                    (noBuffering ? FILE_FLAG_NO_BUFFERING : 0)|FILE_FLAG_SEQUENTIAL_SCAN, NULL); 
    if (fh == INVALID_HANDLE_VALUE) {
        return GetLastError();
    }
    mh = NULL;
    mmapAddr = NULL;
    sharedName = NULL;
    return ok;
}

⌨️ 快捷键说明

复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?