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
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463// 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" "math" "path" "runtime/debug" "sync" "sync/atomic" "time"
"github.com/88250/lute/parse" "github.com/siyuan-note/eventbus" "github.com/siyuan-note/logging" "github.com/siyuan-note/siyuan/kernel/task" "github.com/siyuan-note/siyuan/kernel/treenode" "github.com/siyuan-note/siyuan/kernel/util")
var ( operationQueue []*dbQueueOperation dbQueueLock = sync.Mutex{} dbQueueCond = sync.NewCond(&dbQueueLock) txLock = sync.Mutex{})
type dbQueueOperation struct { inQueueTime time.Time action string // upsert/delete/delete_id/rename/rename_sub_tree/delete_box/delete_box_refs/index/delete_ids/update_block_content/delete_assets indexTree *parse.Tree // index upsertTree *parse.Tree // upsert/update_refs/delete_refs removeTreeBox, removeTreePath string // delete removeTreeID string // delete_id removeTreeIDs []string // delete_ids box string // delete_box/delete_box_refs/index renameTree *parse.Tree // rename/rename_sub_tree block *Block // update_block_content id string // index_node removeAssetHashes []string // delete_assets}
func FlushTxJob() { task.AppendTask(task.DatabaseIndexCommit, FlushQueue)}
func WaitFlushTx() { dbQueueLock.Lock() defer dbQueueLock.Unlock()
var printLog, lastPrintLog bool var i int
for len(operationQueue) > 0 || flushingTx.Load() { if i == 0 { // 第一次等待时使用较短的超时 dbQueueCond.Wait() } else { // 后续等待添加超时检测,用于打印警告日志 timer := time.AfterFunc(50*time.Millisecond, func() { dbQueueCond.Broadcast() }) dbQueueCond.Wait() timer.Stop() }
i++ if 200 < i && !printLog { // 10s 后打日志 logging.LogWarnf("database is writing: \n%s", logging.ShortStack()) printLog = true } if 1200 < i && !lastPrintLog { // 60s 后打日志 logging.LogWarnf("database is still writing") lastPrintLog = true } }}
func ClearQueue() { dbQueueLock.Lock() defer dbQueueLock.Unlock() operationQueue = nil}
var flushingTx = atomic.Bool{}
func FlushQueue() { initDatabaseLock.Lock() defer initDatabaseLock.Unlock()
ops := getOperations() total := len(ops) if 1 > total && !flushingTx.Load() { return }
txLock.Lock() flushingTx.Store(true) defer func() { flushingTx.Store(false) txLock.Unlock() // 通知等待的协程队列已刷新完成 dbQueueCond.Broadcast() }()
start := time.Now()
// logging.LogInfof("flushing database queue, total operations [%d]", total)
// 如果有重命名子树的操作,则统计各路径前缀的块树数量,数量较大的话阻塞整个队列,以便尽可能合并重命名子树的操作 var renameSubTreeOp *dbQueueOperation for _, op := range ops { if "rename_sub_tree" == op.action { renameSubTreeOp = op break } } if nil != renameSubTreeOp { childCount := treenode.CountBlockTreesByPathPrefix(path.Dir(renameSubTreeOp.renameTree.Path)) if 512 < childCount { scale := math.Log(float64(childCount)/512.0+1.0) / math.Log(2.0) secs := 1.0 * scale if secs < 1.0 { secs = 1.0 } if secs > 12.0 { secs = 12.0 } logging.LogInfof("rename sub tree [%s] with large child count [%d], sleep [%.2fs] to wait for more operations", renameSubTreeOp.renameTree.Path, childCount, secs) time.Sleep(time.Duration(secs * float64(time.Second))) } }
context := map[string]interface{}{eventbus.CtxPushMsg: eventbus.CtxPushMsgToStatusBar} if 512 < len(ops) { disableCache() defer enableCache() }
groupOpsTotal := map[string]int{} for _, op := range ops { groupOpsTotal[op.action]++ }
groupOpsCurrent := map[string]int{} for i, op := range ops { if util.IsExiting.Load() { return }
tx, err := beginTx() if err != nil { return }
groupOpsCurrent[op.action]++ context["current"] = groupOpsCurrent[op.action] context["total"] = groupOpsTotal[op.action] if err = execOp(op, tx, context); err != nil { tx.Rollback() closeTxPreparedStmts(tx) logging.LogErrorf("queue operation [%s] failed: %s", op.action, err) continue }
if err = commitTx(tx); err != nil { logging.LogErrorf("commit tx failed: %s", err) continue }
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 op tx [%dms]", elapsed) }
// Push database index commit event https://github.com/siyuan-note/siyuan/issues/8814 util.BroadcastByType("main", "databaseIndexCommit", 0, "", nil)
eventbus.Publish(eventbus.EvtSQLIndexFlushed)}
func execOp(op *dbQueueOperation, tx *sql.Tx, context map[string]interface{}) (err error) { switch op.action { case "index": err = indexTree(tx, op.indexTree, context) case "upsert": err = upsertTree(tx, op.upsertTree, context) case "delete": err = batchDeleteByPathPrefix(tx, op.removeTreeBox, op.removeTreePath) case "delete_id": err = deleteByRootID(tx, op.removeTreeID, context) case "delete_ids": err = batchDeleteByRootIDs(tx, op.removeTreeIDs, context) case "rename": err = batchUpdateHPath(tx, op.renameTree, context) if err != nil { break } err = updateRootContent(tx, path.Base(op.renameTree.HPath), op.renameTree.Root.IALAttr("updated"), op.renameTree.ID) case "rename_sub_tree": err = batchUpdatePath(tx, op.renameTree, context) case "delete_box": err = deleteByBoxTx(tx, op.box) case "delete_box_refs": err = deleteRefsByBoxTx(tx, op.box) case "update_refs": err = upsertRefs(tx, op.upsertTree) case "delete_refs": err = deleteRefs(tx, op.upsertTree) case "update_block_content": err = updateBlockContent(tx, op.block) case "delete_assets": err = deleteAssetsByHashes(tx, op.removeAssetHashes) case "index_node": err = indexNode(tx, op.id) default: msg := fmt.Sprintf("unknown operation [%s]", op.action) logging.LogErrorf(msg) err = errors.New(msg) } return}
func IndexNodeQueue(id string) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{id: id, inQueueTime: time.Now(), action: "index_node"} for i, op := range operationQueue { if "index_node" == op.action && op.id == id { operationQueue[i] = newOp return } } appendOperation(newOp)}
func BatchRemoveAssetsQueue(hashes []string) { if 1 > len(hashes) { return }
dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{removeAssetHashes: hashes, inQueueTime: time.Now(), action: "delete_assets"} appendOperation(newOp)}
func UpdateBlockContentQueue(block *Block) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{block: block, inQueueTime: time.Now(), action: "update_block_content"} for i, op := range operationQueue { if "update_block_content" == op.action && op.block.ID == block.ID { operationQueue[i] = newOp return } } appendOperation(newOp)}
func DeleteRefsTreeQueue(tree *parse.Tree) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{upsertTree: tree, inQueueTime: time.Now(), action: "delete_refs"} for i, op := range operationQueue { if "delete_refs" == op.action && op.upsertTree.ID == tree.ID { operationQueue[i] = newOp return } } appendOperation(newOp)}
func UpdateRefsTreeQueue(tree *parse.Tree) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{upsertTree: tree, inQueueTime: time.Now(), action: "update_refs"} for i, op := range operationQueue { if "update_refs" == op.action && op.upsertTree.ID == tree.ID { operationQueue[i] = newOp return } } appendOperation(newOp)}
func DeleteBoxRefsQueue(boxID string) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{box: boxID, inQueueTime: time.Now(), action: "delete_box_refs"} for i, op := range operationQueue { if "delete_box_refs" == op.action && op.box == boxID { operationQueue[i] = newOp return } } appendOperation(newOp)}
func DeleteBoxQueue(boxID string) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{box: boxID, inQueueTime: time.Now(), action: "delete_box"} for i, op := range operationQueue { if "delete_box" == op.action && op.box == boxID { operationQueue[i] = newOp return } } appendOperation(newOp)}
func IndexTreeQueue(tree *parse.Tree) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{indexTree: tree, inQueueTime: time.Now(), action: "index"} for i, op := range operationQueue { if "index" == op.action && op.indexTree.ID == tree.ID { // 相同树则覆盖 operationQueue[i] = newOp return } } appendOperation(newOp)}
func UpsertTreeQueue(tree *parse.Tree) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{upsertTree: tree, inQueueTime: time.Now(), action: "upsert"} for i, op := range operationQueue { if "upsert" == op.action && op.upsertTree.ID == tree.ID { // 相同树则覆盖 operationQueue[i] = newOp return } } appendOperation(newOp)}
func RenameTreeQueue(tree *parse.Tree) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{ renameTree: tree, inQueueTime: time.Now(), action: "rename", } for i, op := range operationQueue { if "rename" == op.action && op.renameTree.ID == tree.ID { // 相同树则覆盖 operationQueue[i] = newOp return } } appendOperation(newOp)}
func RenameSubTreeQueue(tree *parse.Tree) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{ renameTree: tree, inQueueTime: time.Now(), action: "rename_sub_tree", } for i, op := range operationQueue { if "rename_sub_tree" == op.action && op.renameTree.ID == tree.ID { // 相同树则覆盖 operationQueue[i] = newOp return } } appendOperation(newOp)}
func RemoveTreeQueue(rootID string) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{removeTreeID: rootID, inQueueTime: time.Now(), action: "delete_id"} for i, op := range operationQueue { if "delete_id" == op.action && op.removeTreeID == rootID { operationQueue[i] = newOp return } } appendOperation(newOp)}
func BatchRemoveTreeQueue(rootIDs []string) { if 1 > len(rootIDs) { return }
dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{removeTreeIDs: rootIDs, inQueueTime: time.Now(), action: "delete_ids"} appendOperation(newOp)}
func RemoveTreePathQueue(treeBox, treePathPrefix string) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
newOp := &dbQueueOperation{removeTreeBox: treeBox, removeTreePath: treePathPrefix, inQueueTime: time.Now(), action: "delete"} for i, op := range operationQueue { if "delete" == op.action && (op.removeTreeBox == treeBox && op.removeTreePath == treePathPrefix) { operationQueue[i] = newOp return } } appendOperation(newOp)}
func getOperations() (ops []*dbQueueOperation) { dbQueueLock.Lock() defer dbQueueLock.Unlock()
ops = operationQueue operationQueue = nil return}
func appendOperation(op *dbQueueOperation) { operationQueue = append(operationQueue, op) eventbus.Publish(eventbus.EvtSQLIndexChanged)}