package database import ( "errors" "fmt" "log" "net" "os" "sync" "time" "gordenko.dev/dima/qb" "gordenko.dev/dima/qb/inbox" "gordenko.dev/dima/qb/pagecache" "gordenko.dev/dima/qb/recovery" "gordenko.dev/dima/qb/storage" "gordenko.dev/dima/qb/worker" ) type Database struct { mutex sync.Mutex dir string databaseName string workerInbox *inbox.Inbox worker *worker.Worker storage *storage.Writer indexCache *pagecache.PageCache dataCache *pagecache.PageCache tcpPort int logfile *os.File logger *log.Logger exitCh chan struct{} waitGroup *sync.WaitGroup } type Options struct { TCPPort int Dir string DatabaseName string Logfile *os.File ExitCh chan struct{} WaitGroup *sync.WaitGroup } func New(opt Options) (_ *Database, err error) { if opt.TCPPort <= 0 { return nil, errors.New("TCPPort option is required") } if opt.Dir == "" { return nil, errors.New("Dir option is required") } if opt.DatabaseName == "" { return nil, errors.New("DatabaseName option is required") } if opt.Logfile == nil { return nil, errors.New("Logfile option is required") } if opt.ExitCh == nil { return nil, errors.New("ExitCh option is required") } if opt.WaitGroup == nil { return nil, errors.New("WaitGroup option is required") } s := &Database{ dir: opt.Dir, databaseName: opt.DatabaseName, tcpPort: opt.TCPPort, logfile: opt.Logfile, logger: log.New(opt.Logfile, "", log.LstdFlags), exitCh: opt.ExitCh, waitGroup: opt.WaitGroup, } recoveryReport, err := recovery.Recovery(s.dir, s.databaseName) if err != nil { return nil, fmt.Errorf("recovery.Recovery: %s", err) } //fmt.Printf("%#v\n", recoveryReport) storageInbox := inbox.New() s.workerInbox = inbox.New() // create separately, because needed for page cache dataFile, err := os.OpenFile( qb.GetDataFilePath(opt.Dir, opt.DatabaseName), os.O_CREATE|os.O_RDWR, 0666, ) if err != nil { return nil, err } indexFile, err := os.OpenFile( qb.GetIndexFilePath(opt.Dir, opt.DatabaseName), os.O_CREATE|os.O_RDWR, 0666, ) if err != nil { return nil, err } s.indexCache, err = pagecache.New(pagecache.Options{ File: indexFile, PageSize: storage.IndexPageSize, VerifyPageCRC: storage.VerifyIndexPageCRC32, }) if err != nil { return nil, fmt.Errorf("index pagecache.New: %s", err) } s.dataCache, err = pagecache.New(pagecache.Options{ File: dataFile, PageSize: storage.DataPageSize, VerifyPageCRC: storage.VerifyDataPageCRC32, }) if err != nil { return nil, fmt.Errorf("data pagecache.New: %s", err) } s.storage, err = storage.NewWriter(storage.WriterOptions{ Inbox: storageInbox, WorkerInbox: s.workerInbox, Dir: s.dir, DatabaseName: s.databaseName, WAL: recoveryReport.WAL, IndexFreeList: recoveryReport.IndexFreeList, DataFreeList: recoveryReport.DataFreeList, DataFile: dataFile, IndexFile: indexFile, ExitCh: s.exitCh, WaitGroup: s.waitGroup, }) if err != nil { return nil, fmt.Errorf("storage.NewWriter: %s", err) } s.worker = worker.New(worker.Options{ Inbox: s.workerInbox, StorageInbox: storageInbox, ReplayMetrics: recoveryReport.ReplayMetrics, Dir: opt.Dir, DatabaseName: opt.DatabaseName, ExitCh: opt.ExitCh, WaitGroup: opt.WaitGroup, }) return s, nil } func (s *Database) ListenAndServe() (err error) { listener, err := net.Listen("tcp", fmt.Sprintf(":%d", s.tcpPort)) if err != nil { return fmt.Errorf("net.Listen: %s; port=%d", err, s.tcpPort) } //s.waitGroup.Add(1) go s.storage.Run() s.waitGroup.Add(1) go s.worker.Run() s.logger.Println("database started") for { conn, err := listener.Accept() if err != nil { s.logger.Printf("listener.Accept: %s\n", err) time.Sleep(time.Second) } else { go s.handleTCPConn(conn) } } } type ContinueFullScanReq struct { MetricType qb.MetricType FracDigits byte ResponseWriter qb.AtreeMeasureConsumer LastPageNo uint32 } func (s *Database) ContinueFullScan(req ContinueFullScanReq) error { //fmt.Printf("ContinueFullScan from page %d\n", req.LastPageNo) buf, err := s.dataCache.FetchPage(req.LastPageNo) if err != nil { return fmt.Errorf("dataCache.FetchPage(%d): %s", req.LastPageNo, err) } treeCursor, err := storage.NewBackwardCursor(storage.BackwardCursorOptions{ MetricType: req.MetricType, FracDigits: req.FracDigits, PageNo: req.LastPageNo, PageData: buf, FetchDataPage: s.dataCache.FetchPage, ReleasePage: s.dataCache.ReleasePage, }) if err != nil { return err } defer treeCursor.Close() idx := 0 for { timestamp, value, done, err := treeCursor.Prev() if err != nil { return err } if done { return nil } //fmt.Printf(" %d: %s => %.2f\n", idx, formatTime(timestamp), value) req.ResponseWriter.Feed(timestamp, value) idx++ } } type ContinueRangeScanReq struct { MetricType qb.MetricType FracDigits byte ResponseWriter qb.AtreeMeasureConsumer LastPageNo uint32 Since uint32 } func (s *Database) ContinueRangeScan(req ContinueRangeScanReq) error { buf, err := s.dataCache.FetchPage(req.LastPageNo) if err != nil { return fmt.Errorf("dataCache.FetchPage(%d): %s", req.LastPageNo, err) } treeCursor, err := storage.NewBackwardCursor(storage.BackwardCursorOptions{ MetricType: req.MetricType, FracDigits: req.FracDigits, PageNo: req.LastPageNo, PageData: buf, FetchDataPage: s.dataCache.FetchPage, ReleasePage: s.dataCache.ReleasePage, }) if err != nil { return err } defer treeCursor.Close() for { timestamp, value, done, err := treeCursor.Prev() if err != nil { return err } if done { return nil } req.ResponseWriter.Feed(timestamp, value) if timestamp < req.Since { return nil } } } type RangeScanReq struct { MetricType qb.MetricType FracDigits byte ResponseWriter qb.AtreeMeasureConsumer Since uint32 Until uint32 LastPageNo uint32 IsDataPage bool } func (s *Database) RangeScan(req RangeScanReq) error { var ( pageNo uint32 buf []byte err error ) if req.IsDataPage { pageNo = req.LastPageNo buf, err = s.dataCache.FetchPage(pageNo) if err != nil { return fmt.Errorf("dataCache.FetchPage(%d): %s", pageNo, err) } } else { pageNo, buf, err = s.findDataPage(req.LastPageNo, req.Until) if err != nil { return err } } cursor, err := storage.NewBackwardCursor(storage.BackwardCursorOptions{ MetricType: req.MetricType, FracDigits: req.FracDigits, PageNo: pageNo, PageData: buf, FetchDataPage: s.dataCache.FetchPage, ReleasePage: s.dataCache.ReleasePage, }) if err != nil { return err } defer cursor.Close() for { timestamp, value, done, err := cursor.Prev() if err != nil { return err } if done { return nil } if timestamp <= req.Until { req.ResponseWriter.Feed(timestamp, value) if timestamp < req.Since { // - записи, удовлетворяющие временным рамкам, закончились. return nil } } } } // (records payload only, isZeroLevel, timestamp) func (s *Database) findDataPage(foundPageNo uint32, timestamp uint32) (uint32, []byte, error) { for { buf, err := s.indexCache.FetchPage(foundPageNo) if err != nil { return 0, nil, fmt.Errorf("fetchIndexPage(%d): %s", foundPageNo, err) } toReleaseIndexPageNo := foundPageNo foundPageNo = storage.FindPageOnIndexPage(buf, timestamp) s.indexCache.ReleasePage(toReleaseIndexPageNo) if storage.IsZeroLevelPage(buf) { buf, err := s.dataCache.FetchPage(foundPageNo) if err != nil { return 0, nil, fmt.Errorf("fetchDataPage(%d): %s", foundPageNo, err) } return foundPageNo, buf, nil } } } type PathLeg struct { PageNo uint32 Data []byte } type PathToDataPage struct { Legs []PathLeg LastPageNo uint32 } // func (s *Database) FindPathToLastPage(rootPageNo uint32) (_ PathToDataPage, err error) { // // var ( // // pageNo = rootPageNo // // legs []PathLeg // // ) // // for { // // var buf []byte // // buf, err = s.fetchIndexPage(pageNo) // // if err != nil { // // err = fmt.Errorf("FetchIndexPage(%d): %s", pageNo, err) // // return // // } // // legs = append(legs, PathLeg{ // // PageNo: pageNo, // // Data: buf, // // // childIdx не нужен // // }) // // foundPageNo := getLastPageNo(buf) // // // fix // // if buf[isLastLevelIdx] == 1 { // // return PathToDataPage{ // // Legs: legs, // // LastPageNo: foundPageNo, // // }, nil // // } // // // вглубь // // pageNo = foundPageNo // // } // return // } // DELETE // func (s *Atree) DeletePages(pageNumbers []uint32) { // s.mutex.Lock() // for _, pageNo := range pageNumbers { // delete(s.pages, pageNo) // } // s.mutex.Unlock() // } type Level struct { Idx int PageNumbers []uint32 } func (s *Database) GetAllPages(pageNumbers []uint32) ([]uint32, []uint32, error) { levels := []*Level{ { PageNumbers: pageNumbers, Idx: 0, }, } return s.collectPages(levels) } func (s *Database) collectPages(levels []*Level) (indexPageNumbers []uint32, dataPageNumbers []uint32, err error) { var ( buf []byte pageNumbers []uint32 ) for { if len(levels) == 0 { return } var ( lastIdx = len(levels) - 1 level = levels[lastIdx] pageNo = level.PageNumbers[level.Idx] ) if level.Idx < len(level.PageNumbers) { buf, err = s.indexCache.FetchPage(pageNo) if err != nil { return nil, nil, fmt.Errorf("indexCache.FetchPage(%d): %s", pageNo, err) } pageNumbers = storage.ListPageNumbers(buf) s.indexCache.ReleasePage(pageNo) if storage.IsZeroLevelPage(buf) { dataPageNumbers = append(dataPageNumbers, pageNumbers...) } else { levels = append(levels, &Level{ PageNumbers: pageNumbers, }) indexPageNumbers = append(indexPageNumbers, pageNumbers...) } level.Idx++ } else { levels = levels[:lastIdx] } } } type DeleteSinceReport struct { IndexPageNumbers []uint32 DataPageNumbers []uint32 DataPage []byte IndexLevelTails []storage.IndexLevelTail } // type Level struct { // Idx int // PageNumbers []uint32 // } // fix - find by (since - 1) ? // [since, ...] func (s *Database) DeleteSince(since uint32, indexLevelTail storage.IndexLevelTail) (_ DeleteSinceReport, err error) { var ( indexPageNumbers []uint32 dataPageNumbers []uint32 dataPage []byte buf []byte indexLevelTails []storage.IndexLevelTail foundPageNo uint32 ) result := storage.DeleteSinceOnIndexTail(indexLevelTail, since) // indexLevelTails = append(indexLevelTails, storage.IndexLevelTail{ Buffer: indexLevelTail.Buffer, RecordsCount: result.RecordsCount, }) levels := []*Level{ { PageNumbers: result.PageNumbers, Idx: 1, }, } foundPageNo = result.PageNumbers[0] for { buf, err = s.indexCache.FetchPage(foundPageNo) if err != nil { err = fmt.Errorf("indexCache.FetchPage(%d): %s", foundPageNo, err) return } toReleaseIndexPageNo := foundPageNo result := storage.DeleteSinceOnIndexPage(buf, since) s.indexCache.ReleasePage(toReleaseIndexPageNo) // foundPageNo = result.PageNumbers[0] levels = append(levels, &Level{ PageNumbers: result.PageNumbers, Idx: 1, // 1st - found pageNo }) indexLevelTails = append(indexLevelTails, storage.IndexLevelTail{ Buffer: result.Buffer, RecordsCount: result.RecordsCount, }) if storage.IsZeroLevelPage(buf) { dataPage, err = s.dataCache.FetchPage(foundPageNo) if err != nil { err = fmt.Errorf("dataCache.FetchPage(%d): %s", foundPageNo, err) return } dataPageNumbers = append(dataPageNumbers, result.PageNumbers...) break } else { indexPageNumbers = append(indexPageNumbers, result.PageNumbers...) } } indexPageNumbers, dataPageNumbers, err = s.collectPages(levels) if err != nil { return } return DeleteSinceReport{ IndexPageNumbers: indexPageNumbers, DataPageNumbers: dataPageNumbers, DataPage: dataPage, IndexLevelTails: indexLevelTails, }, nil } // [..., until] func (s *Database) DeleteUntil(since uint32, indexLevelTail storage.IndexLevelTail) (_ DeleteSinceReport, err error) { var ( indexPageNumbers []uint32 dataPageNumbers []uint32 dataPage []byte buf []byte indexLevelTails []storage.IndexLevelTail foundPageNo uint32 ) result := storage.DeleteSinceOnIndexTail(indexLevelTail, since) // indexLevelTails = append(indexLevelTails, storage.IndexLevelTail{ Buffer: indexLevelTail.Buffer, RecordsCount: result.RecordsCount, }) levels := []*Level{ { PageNumbers: result.PageNumbers, Idx: 1, }, } foundPageNo = result.PageNumbers[0] for { buf, err = s.indexCache.FetchPage(foundPageNo) if err != nil { err = fmt.Errorf("indexCache.FetchPage(%d): %s", foundPageNo, err) return } toReleaseIndexPageNo := foundPageNo result := storage.DeleteSinceOnIndexPage(buf, since) s.indexCache.ReleasePage(toReleaseIndexPageNo) // foundPageNo = result.PageNumbers[0] levels = append(levels, &Level{ PageNumbers: result.PageNumbers, Idx: 1, // 1st - found pageNo }) indexLevelTails = append(indexLevelTails, storage.IndexLevelTail{ Buffer: result.Buffer, RecordsCount: result.RecordsCount, }) if storage.IsZeroLevelPage(buf) { dataPage, err = s.dataCache.FetchPage(foundPageNo) if err != nil { err = fmt.Errorf("dataCache.FetchPage(%d): %s", foundPageNo, err) return } dataPageNumbers = append(dataPageNumbers, result.PageNumbers...) break } else { indexPageNumbers = append(indexPageNumbers, result.PageNumbers...) } } indexPageNumbers, dataPageNumbers, err = s.collectPages(levels) if err != nil { return } return DeleteSinceReport{ IndexPageNumbers: indexPageNumbers, DataPageNumbers: dataPageNumbers, DataPage: dataPage, IndexLevelTails: indexLevelTails, }, nil } const datetimeLayout = "2006-01-02 15:04:05" func formatTime(timestamp uint32) string { tm := time.Unix(int64(timestamp), 0) return tm.Format(datetimeLayout) }