Something went wrong. Try again.
A privacy-first, self-hosted, fully open source personal knowledge management software, written in typescript and golang. (PERSONAL FORK)
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164// SiYuan - Refactor your thinking// Copyright (c) 2020-present, b3log.org//// This program is free software: you can redistribute it and/or modify// it under the terms of the GNU Affero General Public License as published by// the Free Software Foundation, either version 3 of the License, or// (at your option) any later version.//// This program is distributed in the hope that it will be useful,// but WITHOUT ANY WARRANTY; without even the implied warranty of// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the// GNU Affero General Public License for more details.//// You should have received a copy of the GNU Affero General Public License// along with this program. If not, see <https://www.gnu.org/licenses/>.
package sql
import ( "database/sql" "errors" "fmt" "os" "path/filepath" "runtime/debug" "strings" "sync" "time"
"github.com/siyuan-note/eventbus" "github.com/siyuan-note/logging" "github.com/siyuan-note/siyuan/kernel/task" "github.com/siyuan-note/siyuan/kernel/util")
var ( historyOperationQueue []*historyDBQueueOperation historyDBQueueLock = sync.Mutex{} historyTxLock = sync.Mutex{})
type historyDBQueueOperation struct { inQueueTime time.Time action string // index/deleteOutdated
histories []*History // index before int64 // deleteOutdated}
func FlushHistoryTxJob() { task.AppendTask(task.HistoryDatabaseIndexCommit, FlushHistoryQueue)}
func FlushHistoryQueue() { ops := getHistoryOperations() total := len(ops) if 1 > total { return }
historyTxLock.Lock() defer historyTxLock.Unlock() start := time.Now()
groupOpsTotal := map[string]int{} for _, op := range ops { groupOpsTotal[op.action]++ }
context := map[string]interface{}{eventbus.CtxPushMsg: eventbus.CtxPushMsgToStatusBar} groupOpsCurrent := map[string]int{} for i, op := range ops { if util.IsExiting.Load() { return }
tx, err := beginHistoryTx() if err != nil { return }
groupOpsCurrent[op.action]++ context["current"] = groupOpsCurrent[op.action] context["total"] = groupOpsTotal[op.action]
if err = execHistoryOp(op, tx, context); err != nil { tx.Rollback() logging.LogErrorf("queue operation failed: %s", err)
if 0 < len(op.histories) { dir := op.histories[0].Path[:strings.Index(op.histories[0].Path, "/")] dirPath := filepath.Join(util.HistoryDir, dir) if removeErr := os.RemoveAll(dirPath); nil != removeErr { logging.LogErrorf("remove corrupted history dir [%s] failed: %s", dirPath, removeErr) } }
eventbus.Publish(util.EvtSQLHistoryRebuild) return }
if err = commitHistoryTx(tx); err != nil { logging.LogErrorf("commit tx failed: %s", err) return }
if 16 < i && 0 == i%128 { debug.FreeOSMemory() } }
if 128 < total { debug.FreeOSMemory() }
elapsed := time.Now().Sub(start).Milliseconds() if 7000 < elapsed { logging.LogInfof("database history op tx [%dms]", elapsed) }}
func execHistoryOp(op *historyDBQueueOperation, tx *sql.Tx, context map[string]interface{}) (err error) { switch op.action { case "index": err = insertHistories(tx, op.histories, context) case "deleteOutdated": err = deleteOutdatedHistories(tx, op.before, context) default: msg := fmt.Sprintf("unknown history operation [%s]", op.action) logging.LogErrorf(msg) err = errors.New(msg) } return}
func DeleteOutdatedHistories(before int64) { historyDBQueueLock.Lock() defer historyDBQueueLock.Unlock()
newOp := &historyDBQueueOperation{inQueueTime: time.Now(), action: "deleteOutdated", before: before} historyOperationQueue = append(historyOperationQueue, newOp)}
func IndexHistoriesQueue(histories []*History) { if 1 > len(histories) { return }
historyDBQueueLock.Lock() defer historyDBQueueLock.Unlock()
newOp := &historyDBQueueOperation{inQueueTime: time.Now(), action: "index", histories: histories} historyOperationQueue = append(historyOperationQueue, newOp)}
func getHistoryOperations() (ops []*historyDBQueueOperation) { historyDBQueueLock.Lock() defer historyDBQueueLock.Unlock()
ops = historyOperationQueue historyOperationQueue = nil return}