/
githubmirror
/
client
Обзор
Документация
Войти
/
githubmirror
/
client
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
go/chat/storage/storage_blockengine.go
480 строк
14 KB
zoom-ua
enable gosec (#28775)
08 янв 2026, 19:26
Не верифицирован
08 янв 2026, 19:26
5670f73
Код
Авторство
О чём код?
package storage import ( "context" "fmt" "github.com/keybase/client/go/chat/globals" "github.com/keybase/client/go/chat/utils" "github.com/keybase/client/go/libkb" "github.com/keybase/client/go/protocol/chat1" "github.com/keybase/client/go/protocol/gregor1" "golang.org/x/crypto/nacl/secretbox" ) const ( blockIndexVersion = 8 blockSize = 100 ) type blockEngine struct { globals.Contextified utils.DebugLabeler } func newBlockEngine(g *globals.Context) *blockEngine { return &blockEngine{ Contextified: globals.NewContextified(g), DebugLabeler: utils.NewDebugLabeler(g.ExternalG(), "BlockEngine", true), } } type blockIndex struct { Version int ServerVersion int ConvID chat1.ConversationID UID gregor1.UID MaxBlock int BlockSize int } type block struct { BlockID int Msgs [blockSize]chat1.MessageUnboxed } type boxedBlock struct { V int N [24]byte E []byte } func (be *blockEngine) makeBlockKey(convID chat1.ConversationID, uid gregor1.UID, blockID int) libkb.DbKey { return libkb.DbKey{ Typ: libkb.DBChatBlocks, Key: fmt.Sprintf("bl:%s:%s:%d", uid, convID, blockID), } } func (be *blockEngine) getBlockNumber(id chat1.MessageID) int { return int(id) / blockSize //nolint:gosec // G115: MessageID to block number calculation, safe to convert } func (be *blockEngine) getBlockPosition(id chat1.MessageID) int { return int(id) % blockSize //nolint:gosec // G115: MessageID to block position calculation, safe to convert } func (be *blockEngine) getMsgID(blockNum, blockPos int) chat1.MessageID { return chat1.MessageID(blockNum*blockSize + blockPos) //nolint:gosec // G115: Block arithmetic to MessageID, safe to convert } func (be *blockEngine) createBlockIndex(ctx context.Context, key libkb.DbKey, convID chat1.ConversationID, uid gregor1.UID, ) (bi blockIndex, err Error) { be.Debug(ctx, "createBlockIndex: creating new block index: convID: %s uid: %s", convID, uid) // Grab latest server version to tag local data with srvVers, serr := be.G().ServerCacheVersions.Fetch(ctx) if serr != nil { return blockIndex{}, NewInternalError(ctx, be.DebugLabeler, "createBlockIndex: failed to get server versions: %s", serr.Error()) } bi = blockIndex{ Version: blockIndexVersion, ServerVersion: srvVers.BodiesVers, ConvID: convID, UID: uid, MaxBlock: -1, BlockSize: blockSize, } dat, ierr := encode(bi) if ierr != nil { return bi, NewInternalError(ctx, be.DebugLabeler, "createBlockIndex: failed to encode %s", ierr) } if ierr = be.G().LocalChatDb.PutRaw(key, dat); ierr != nil { return bi, NewInternalError(ctx, be.DebugLabeler, "createBlockIndex: failed to write: %s", ierr) } return bi, nil } func (be *blockEngine) readBlockIndex(ctx context.Context, convID chat1.ConversationID, uid gregor1.UID) (blockIndex, Error) { key := makeBlockIndexKey(convID, uid) raw, found, err := be.G().LocalChatDb.GetRaw(key) if err != nil { return blockIndex{}, NewInternalError(ctx, be.DebugLabeler, "readBlockIndex: failed to read index block: %s", err.Error()) } if !found { // If not found, create a new one and return it be.Debug(ctx, "readBlockIndex: no block index found, creating: convID: %d uid: %s", convID, uid) return be.createBlockIndex(ctx, key, convID, uid) } // Decode and return var bi blockIndex if err = decode(raw, &bi); err != nil { return bi, NewInternalError(ctx, be.DebugLabeler, "readBlockIndex: failed to decode: %s", err.Error()) } if bi.Version != blockIndexVersion { be.Debug(ctx, "readBlockInbox: version mismatch, creating new index") return be.createBlockIndex(ctx, key, convID, uid) } // Check server version if _, err = be.G().ServerCacheVersions.MatchBodies(ctx, bi.ServerVersion); err != nil { be.Debug(ctx, "readBlockInbox: server version error: %s, creating new index", err.Error()) return be.createBlockIndex(ctx, key, convID, uid) } return bi, nil } type bekey string var ( bebikey bekey = "bebi" beskkey bekey = "besk" ) func (be *blockEngine) Init(ctx context.Context, key [32]byte, convID chat1.ConversationID, uid gregor1.UID, ) (context.Context, Error) { ctx = context.WithValue(ctx, beskkey, key) bi, err := be.readBlockIndex(ctx, convID, uid) if err != nil { return ctx, err } ctx = context.WithValue(ctx, bebikey, &bi) return ctx, nil } func (be *blockEngine) fetchBlockIndex(ctx context.Context, convID chat1.ConversationID, uid gregor1.UID, ) (bi blockIndex, err Error) { var ok bool val := ctx.Value(bebikey) if bi, ok = val.(blockIndex); !ok { bi, err = be.readBlockIndex(ctx, convID, uid) if err != nil { return bi, err } } be.Debug(ctx, "fetchBlockIndex: maxBlock: %d", bi.MaxBlock) return bi, err } func (be *blockEngine) fetchSecretKey(ctx context.Context) (key [32]byte, err Error) { var ok bool val := ctx.Value(beskkey) if key, ok = val.([32]byte); !ok { return key, MiscError{Msg: "secret key not in context"} } return key, nil } func (be *blockEngine) createBlockSingle(ctx context.Context, bi blockIndex, blockID int) (block, Error) { be.Debug(ctx, "createBlockSingle: creating block: %d", blockID) // Write out new block b := block{BlockID: blockID} if cerr := be.writeBlock(ctx, bi, b); cerr != nil { return block{}, NewInternalError(ctx, be.DebugLabeler, "createBlockSingle: failed to write block: %s", cerr.Message()) } return b, nil } func (be *blockEngine) createBlock(ctx context.Context, bi *blockIndex, blockID int) (block, Error) { // Create all the blocks up to the one we want var b block for i := bi.MaxBlock + 1; i <= blockID; i++ { b, err := be.createBlockSingle(ctx, *bi, i) if err != nil { return b, err } } // Update block index with new block bi.MaxBlock = blockID dat, err := encode(bi) if err != nil { return block{}, NewInternalError(ctx, be.DebugLabeler, "createBlock: failed to encode block: %s", err.Error()) } err = be.G().LocalChatDb.PutRaw(makeBlockIndexKey(bi.ConvID, bi.UID), dat) if err != nil { return block{}, NewInternalError(ctx, be.DebugLabeler, "createBlock: failed to write index: %s", err.Error()) } return b, nil } func (be *blockEngine) getBlock(ctx context.Context, bi blockIndex, id chat1.MessageID) (block, Error) { if id == 0 { return block{}, NewInternalError(ctx, be.DebugLabeler, "getBlock: invalid block id: %d", id) } bn := be.getBlockNumber(id) if bn > bi.MaxBlock { be.Debug(ctx, "getBlock(): missed high: id: %d maxblock: %d", bn, bi.MaxBlock) return block{}, MissError{} } return be.readBlock(ctx, bi, bn) } func (be *blockEngine) readBlock(ctx context.Context, bi blockIndex, id int) (res block, err Error) { be.Debug(ctx, "readBlock: reading block: %d", id) // Manage in memory cache if b, ok := blockEngineMemCache.getBlock(ctx, bi.UID, bi.ConvID, id); ok { be.Debug(ctx, "readBlock: cache hit") return b, nil } defer func() { if err == nil { blockEngineMemCache.writeBlock(ctx, bi.UID, bi.ConvID, res) } }() key := be.makeBlockKey(bi.ConvID, bi.UID, id) raw, found, ierr := be.G().LocalChatDb.GetRaw(key) if ierr != nil { return res, NewInternalError(ctx, be.DebugLabeler, "readBlock: failed to read raw: %s", ierr.Error()) } if !found { // Didn't find it for some reason return res, NewInternalError(ctx, be.DebugLabeler, "readBlock: block not found: id: %d", id) } // Decode boxed block var b boxedBlock if ierr := decode(raw, &b); ierr != nil { return res, NewInternalError(ctx, be.DebugLabeler, "readBlock: failed to decode: %s", ierr.Error()) } if b.V > cryptoVersion { return res, NewInternalError(ctx, be.DebugLabeler, "readBlock: bad crypto version: %d current: %d id: %d", b.V, cryptoVersion, id) } // Decrypt block fkey, cerr := be.fetchSecretKey(ctx) if cerr != nil { return res, cerr } pt, ok := secretbox.Open(nil, b.E, &b.N, &fkey) if !ok { return res, NewInternalError(ctx, be.DebugLabeler, "readBlock: failed to decrypt block: %d", id) } // Decode payload if ierr = decode(pt, &res); ierr != nil { return res, NewInternalError(ctx, be.DebugLabeler, "readBlock: failed to decode: %s", ierr.Error()) } return res, nil } func (be *blockEngine) writeBlock(ctx context.Context, bi blockIndex, b block) (err Error) { be.Debug(ctx, "writeBlock: writing out block: %d", b.BlockID) defer func() { if err == nil { blockEngineMemCache.writeBlock(ctx, bi.UID, bi.ConvID, b) } }() // Encode block dat, ierr := encode(b) if ierr != nil { return NewInternalError(ctx, be.DebugLabeler, "writeBlock: failed to encode: %s", ierr.Error()) } // Encrypt block key, cerr := be.fetchSecretKey(ctx) if cerr != nil { return cerr } var nonce []byte nonce, ierr = libkb.RandBytes(24) if ierr != nil { return MiscError{Msg: fmt.Sprintf("encryptMessage: failure to generate nonce: %s", ierr.Error())} } var fnonce [24]byte copy(fnonce[:], nonce) sealed := secretbox.Seal(nil, dat, &fnonce, &key) // Encode encrypted block payload := boxedBlock{ V: cryptoVersion, N: fnonce, E: sealed, } bpayload, ierr := encode(payload) if ierr != nil { return NewInternalError(ctx, be.DebugLabeler, "writeBlock: failed to encode: %s", ierr.Error()) } // Write out encrypted block if ierr := be.G().LocalChatDb.PutRaw(be.makeBlockKey(bi.ConvID, bi.UID, b.BlockID), bpayload); ierr != nil { return NewInternalError(ctx, be.DebugLabeler, "writeBlock: failed to write: %s", ierr.Error()) } return nil } func (be *blockEngine) WriteMessages(ctx context.Context, convID chat1.ConversationID, uid gregor1.UID, msgs []chat1.MessageUnboxed, ) Error { msgIDs := make([]chat1.MessageID, len(msgs)) msgMap := make(map[chat1.MessageID]chat1.MessageUnboxed) for index, msg := range msgs { msgMap[msg.GetMessageID()] = msg msgIDs[index] = msg.GetMessageID() } return be.writeMessagesIDMap(ctx, convID, uid, msgIDs, msgMap) } func (be *blockEngine) writeMessagesIDMap(ctx context.Context, convID chat1.ConversationID, uid gregor1.UID, msgIDs []chat1.MessageID, msgMap map[chat1.MessageID]chat1.MessageUnboxed, ) Error { var err Error var maxB block var newBlock block var lastWritten int docreate := false // Get block index bi, err := be.fetchBlockIndex(ctx, convID, uid) if err != nil { return err } // Sanity check if len(msgIDs) == 0 { return nil } // Get the maximum block (create it if we need to) maxID := msgIDs[0] be.Debug(ctx, "writeMessages: maxID: %d num: %d", maxID, len(msgIDs)) if maxB, err = be.getBlock(ctx, bi, maxID); err != nil { if _, ok := err.(MissError); !ok { return err } docreate = true } if docreate { newBlockID := be.getBlockNumber(maxID) be.Debug(ctx, "writeMessages: block not found (creating): maxID: %d id: %d", maxID, newBlockID) if _, err = be.createBlock(ctx, &bi, newBlockID); err != nil { return NewInternalError(ctx, be.DebugLabeler, "writeMessages: failed to create block: %s", err.Message()) } if maxB, err = be.getBlock(ctx, bi, maxID); err != nil { return NewInternalError(ctx, be.DebugLabeler, "writeMessages: failed to read newly created block: %s", err.Message()) } } // Append to the block newBlock = maxB for index, msgID := range msgIDs { if be.getBlockNumber(msgID) != newBlock.BlockID { be.Debug(ctx, "writeMessages: crossed block boundary, aborting and writing out: msgID: %d", msgID) break } newBlock.Msgs[be.getBlockPosition(msgID)] = msgMap[msgID] lastWritten = index } // Write the block if err = be.writeBlock(ctx, bi, newBlock); err != nil { return NewInternalError(ctx, be.DebugLabeler, "writeMessages: failed to write block: %s", err.Message()) } // We didn't write everything out in this block, move to another one if lastWritten < len(msgIDs)-1 { return be.writeMessagesIDMap(ctx, convID, uid, msgIDs[lastWritten+1:], msgMap) } return nil } func (be *blockEngine) ReadMessages(ctx context.Context, res ResultCollector, convID chat1.ConversationID, uid gregor1.UID, maxID, minID chat1.MessageID, ) (err Error) { // Run all errors through resultCollector defer func() { if err != nil { err = res.Error(err) } }() // Get block index bi, err := be.fetchBlockIndex(ctx, convID, uid) if err != nil { return err } // Get the current block where max ID is found b, err := be.getBlock(ctx, bi, maxID) if err != nil { return err } // Add messages to result set var lastAdded chat1.MessageID maxPos := be.getBlockPosition(maxID) be.Debug(ctx, "readMessages: BID: %d maxPos: %d maxID: %d rc: %s", b.BlockID, maxPos, maxID, res) for index := maxPos; !res.Done() && index >= 0; index-- { if b.BlockID == 0 && index == 0 { // Short circuit out of here if we are on the null message break } msg := b.Msgs[index] // If we have a versioning error but our client now understands the new // version, don't return the error message if msg.GetMessageID() == 0 || (msg.IsError() && msg.Error().ParseableVersion()) { if res.PushPlaceholder(be.getMsgID(b.BlockID, index)) { // If the result collector is happy to receive this blank entry, then don't complain // and proceed as if this was a hit lastAdded = be.getMsgID(b.BlockID, index) be.Debug(ctx, "readMessages: adding placeholder: %d (blockid: %d pos: %d)", lastAdded, b.BlockID, index) continue } be.Debug(ctx, "readMessages: cache entry empty: index: %d block: %d msgID: %d", index, b.BlockID, be.getMsgID(b.BlockID, index)) return MissError{} } else if msg.GetMessageID() <= minID { // If we drop below the min ID, just bail out of here with no error return nil } bMsgID := msg.GetMessageID() // Sanity check if bMsgID != be.getMsgID(b.BlockID, index) { return NewInternalError(ctx, be.DebugLabeler, "chat entry corruption: bMsgID: %d != %d (block: %d pos: %d)", bMsgID, be.getMsgID(b.BlockID, index), b.BlockID, index) } be.Debug(ctx, "readMessages: adding msg_id: %d (blockid: %d pos: %d)", msg.GetMessageID(), b.BlockID, index) lastAdded = msg.GetMessageID() res.Push(msg) } // Check if we read anything, otherwise move to another block and try // again. We check if lastAdded > 0 to avoid overflowing chat1.MessageID // which is a uint type if !res.Done() && b.BlockID > 0 && lastAdded > 0 { return be.ReadMessages(ctx, res, convID, uid, lastAdded-1, minID) } return nil } func (be *blockEngine) ClearMessages(ctx context.Context, convID chat1.ConversationID, uid gregor1.UID, msgIDs []chat1.MessageID, ) Error { msgMap := make(map[chat1.MessageID]chat1.MessageUnboxed) for _, msgID := range msgIDs { msgMap[msgID] = chat1.MessageUnboxed{} } return be.writeMessagesIDMap(ctx, convID, uid, msgIDs, msgMap) }