diff --git a/carstore/bs.go b/carstore/bs.go index 1f86b661..894f41c5 100644 --- a/carstore/bs.go +++ b/carstore/bs.go @@ -4,14 +4,11 @@ import ( "bufio" "bytes" "context" - "encoding/binary" "fmt" "io" "os" "path/filepath" "sort" - "strconv" - "strings" "sync" "sync/atomic" "time" @@ -33,7 +30,6 @@ import ( cbg "github.com/whyrusleeping/cbor-gen" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" - "gorm.io/driver/postgres" "gorm.io/gorm" ) @@ -53,7 +49,7 @@ const MaxSliceLength = 2 << 20 const BigShardThreshold = 2 << 20 type CarStore struct { - meta *gorm.DB + meta *CarStoreGormMeta rootDir string lscLk sync.Mutex @@ -78,79 +74,12 @@ func NewCarStore(meta *gorm.DB, root string) (*CarStore, error) { } return &CarStore{ - meta: meta, + meta: &CarStoreGormMeta{meta: meta}, rootDir: root, lastShardCache: make(map[models.Uid]*CarShard), }, nil } -type UserInfo struct { - gorm.Model - Head string -} - -type CarShard struct { - ID uint `gorm:"primarykey"` - CreatedAt time.Time - - Root models.DbCID `gorm:"index"` - DataStart int64 - Seq int `gorm:"index:idx_car_shards_seq;index:idx_car_shards_usr_seq,priority:2,sort:desc"` - Path string - Usr models.Uid `gorm:"index:idx_car_shards_usr;index:idx_car_shards_usr_seq,priority:1"` - Rev string -} - -type blockRef struct { - ID uint `gorm:"primarykey"` - Cid models.DbCID `gorm:"index"` - Shard uint `gorm:"index"` - Offset int64 - //User uint `gorm:"index"` -} - -type staleRef struct { - ID uint `gorm:"primarykey"` - Cid *models.DbCID - Cids []byte - Usr models.Uid `gorm:"index"` -} - -func (sr *staleRef) getCids() ([]cid.Cid, error) { - if sr.Cid != nil { - return []cid.Cid{sr.Cid.CID}, nil - } - - return unpackCids(sr.Cids) -} - -func packCids(cids []cid.Cid) []byte { - buf := new(bytes.Buffer) - for _, c := range cids { - buf.Write(c.Bytes()) - } - - return buf.Bytes() -} - -func unpackCids(b []byte) ([]cid.Cid, error) { - br := bytes.NewReader(b) - var out []cid.Cid - for { - _, c, err := cid.CidFromReader(br) - if err != nil { - if err == io.EOF { - break - } - return nil, err - } - - out = append(out, c) - } - - return out, nil -} - type userView struct { cs *CarStore user models.Uid @@ -166,17 +95,7 @@ func (uv *userView) HashOnRead(hor bool) { } func (uv *userView) Has(ctx context.Context, k cid.Cid) (bool, error) { - var count int64 - if err := uv.cs.meta. - Model(blockRef{}). - Select("path, block_refs.offset"). - Joins("left join car_shards on block_refs.shard = car_shards.id"). - Where("usr = ? AND cid = ?", uv.user, models.DbCID{CID: k}). - Count(&count).Error; err != nil { - return false, err - } - - return count > 0, nil + return uv.cs.meta.HasUidCid(ctx, uv.user, k) } var CacheHits int64 @@ -197,30 +116,16 @@ func (uv *userView) Get(ctx context.Context, k cid.Cid) (blockformat.Block, erro } atomic.AddInt64(&CacheMiss, 1) - // TODO: for now, im using a join to ensure we only query blocks from the - // correct user. maybe it makes sense to put the user in the blockRef - // directly? tradeoff of time vs space - var info struct { - Path string - Offset int64 - Usr models.Uid - } - if err := uv.cs.meta.Raw(`SELECT - (select path from car_shards where id = block_refs.shard) as path, - block_refs.offset, - (select usr from car_shards where id = block_refs.shard) as usr -FROM block_refs -WHERE - block_refs.cid = ? -LIMIT 1;`, models.DbCID{CID: k}).Scan(&info).Error; err != nil { + path, offset, user, err := uv.cs.meta.LookupBlockRef(ctx, k) + if err != nil { return nil, err } - if info.Path == "" { + if path == "" { return nil, ipld.ErrNotFound{Cid: k} } prefetch := uv.prefetch - if info.Usr != uv.user { + if user != uv.user { blockGetTotalCounterUsrskip.Add(1) prefetch = false } else { @@ -228,9 +133,9 @@ LIMIT 1;`, models.DbCID{CID: k}).Scan(&info).Error; err != nil { } if prefetch { - return uv.prefetchRead(ctx, k, info.Path, info.Offset) + return uv.prefetchRead(ctx, k, path, offset) } else { - return uv.singleRead(ctx, k, info.Path, info.Offset) + return uv.singleRead(ctx, k, path, offset) } } @@ -390,18 +295,13 @@ func (cs *CarStore) getLastShard(ctx context.Context, user models.Uid) (*CarShar return maybeLs, nil } - var lastShard CarShard - // this is often slow (which is why we're caching it) but could be sped up with an extra index: - // CREATE INDEX idx_car_shards_usr_id ON car_shards (usr, seq DESC); - if err := cs.meta.WithContext(ctx).Model(CarShard{}).Limit(1).Order("seq desc").Find(&lastShard, "usr = ?", user).Error; err != nil { - //if err := cs.meta.Model(CarShard{}).Where("user = ?", user).Last(&lastShard).Error; err != nil { - //if err != gorm.ErrRecordNotFound { + lastShard, err := cs.meta.GetLastShard(ctx, user) + if err != nil { return nil, err - //} } - cs.putLastShardCache(&lastShard) - return &lastShard, nil + cs.putLastShardCache(lastShard) + return lastShard, nil } var ErrRepoBaseMismatch = fmt.Errorf("attempted a delta session on top of the wrong previous head") @@ -452,24 +352,27 @@ func (cs *CarStore) ReadOnlySession(user models.Uid) (*DeltaSession, error) { }, nil } +// TODO: incremental is only ever called true, remove the param func (cs *CarStore) ReadUserCar(ctx context.Context, user models.Uid, sinceRev string, incremental bool, w io.Writer) error { ctx, span := otel.Tracer("carstore").Start(ctx, "ReadUserCar") defer span.End() var earlySeq int if sinceRev != "" { - var untilShard CarShard - if err := cs.meta.Where("rev >= ? AND usr = ?", sinceRev, user).Order("rev").First(&untilShard).Error; err != nil { - return fmt.Errorf("finding early shard: %w", err) + var err error + earlySeq, err = cs.meta.SeqForRev(ctx, user, sinceRev) + if err != nil { + return err } - earlySeq = untilShard.Seq } - var shards []CarShard - if err := cs.meta.Order("seq desc").Where("usr = ? AND seq >= ?", user, earlySeq).Find(&shards).Error; err != nil { + // TODO: Why does ReadUserCar want shards seq DESC but CompactUserShards wants seq ASC ? + shards, err := cs.meta.GetUserShardsDesc(ctx, user, earlySeq) + if err != nil { return err } + // TODO: incremental is only ever called true, so this is fine and we can remove the error check if !incremental && earlySeq > 0 { // have to do it the ugly way return fmt.Errorf("nyi") @@ -496,6 +399,8 @@ func (cs *CarStore) ReadUserCar(ctx context.Context, user models.Uid, sinceRev s return nil } +// inner loop part of ReadUserCar +// copy shard blocks from disk to Writer func (cs *CarStore) writeShardBlocks(ctx context.Context, sh *CarShard, w io.Writer) error { ctx, span := otel.Tracer("carstore").Start(ctx, "writeShardBlocks") defer span.End() @@ -519,31 +424,7 @@ func (cs *CarStore) writeShardBlocks(ctx context.Context, sh *CarShard, w io.Wri return nil } -func (cs *CarStore) writeBlockFromShard(ctx context.Context, sh *CarShard, w io.Writer, c cid.Cid) error { - fi, err := os.Open(sh.Path) - if err != nil { - return err - } - defer fi.Close() - - rr, err := car.NewCarReader(fi) - if err != nil { - return err - } - - for { - blk, err := rr.Next() - if err != nil { - return err - } - - if blk.Cid() == c { - _, err := LdWrite(w, c.Bytes(), blk.RawData()) - return err - } - } -} - +// inner loop part of compactBucket func (cs *CarStore) iterateShardBlocks(ctx context.Context, sh *CarShard, cb func(blk blockformat.Block) error) error { fi, err := os.Open(sh.Path) if err != nil { @@ -769,123 +650,18 @@ func (cs *CarStore) putShard(ctx context.Context, shard *CarShard, brefs []map[s ctx, span := otel.Tracer("carstore").Start(ctx, "putShard") defer span.End() - // TODO: there should be a way to create the shard and block_refs that - // reference it in the same query, would save a lot of time - tx := cs.meta.WithContext(ctx).Begin() - - if err := tx.WithContext(ctx).Create(shard).Error; err != nil { - return fmt.Errorf("failed to create shard in DB tx: %w", err) - } - - if !nocache { - cs.putLastShardCache(shard) - } - - for _, ref := range brefs { - ref["shard"] = shard.ID - } - - if err := createBlockRefs(ctx, tx, brefs); err != nil { - return fmt.Errorf("failed to create block refs: %w", err) - } - - if len(rmcids) > 0 { - cids := make([]cid.Cid, 0, len(rmcids)) - for c := range rmcids { - cids = append(cids, c) - } - - if err := tx.Create(&staleRef{ - Cids: packCids(cids), - Usr: shard.Usr, - }).Error; err != nil { - return err - } - } - - err := tx.WithContext(ctx).Commit().Error + err := cs.meta.PutShardAndRefs(ctx, shard, brefs, rmcids) if err != nil { - return fmt.Errorf("failed to commit shard DB transaction: %w", err) - } - - return nil -} - -func createBlockRefs(ctx context.Context, tx *gorm.DB, brefs []map[string]any) error { - ctx, span := otel.Tracer("carstore").Start(ctx, "createBlockRefs") - defer span.End() - - if err := createInBatches(ctx, tx, brefs, 2000); err != nil { return err } - return nil -} - -func generateInsertQuery(data []map[string]any) (string, []any) { - placeholders := strings.Repeat("(?, ?, ?),", len(data)) - placeholders = placeholders[:len(placeholders)-1] // trim trailing comma - - query := "INSERT INTO block_refs (\"cid\", \"offset\", \"shard\") VALUES " + placeholders - - values := make([]any, 0, 3*len(data)) - for _, entry := range data { - values = append(values, entry["cid"], entry["offset"], entry["shard"]) + if !nocache { + cs.putLastShardCache(shard) } - return query, values -} - -// Function to create in batches -func createInBatches(ctx context.Context, tx *gorm.DB, data []map[string]any, batchSize int) error { - for i := 0; i < len(data); i += batchSize { - batch := data[i:] - if len(batch) > batchSize { - batch = batch[:batchSize] - } - - query, values := generateInsertQuery(batch) - - if err := tx.WithContext(ctx).Exec(query, values...).Error; err != nil { - return err - } - } return nil } -func LdWrite(w io.Writer, d ...[]byte) (int64, error) { - var sum uint64 - for _, s := range d { - sum += uint64(len(s)) - } - - buf := make([]byte, 8) - n := binary.PutUvarint(buf, sum) - nw, err := w.Write(buf[:n]) - if err != nil { - return 0, err - } - - for _, s := range d { - onw, err := w.Write(s) - if err != nil { - return int64(nw), err - } - nw += onw - } - - return int64(nw), nil -} - -func setToSlice(s map[cid.Cid]bool) []cid.Cid { - out := make([]cid.Cid, 0, len(s)) - for c := range s { - out = append(out, c) - } - - return out -} - func BlockDiff(ctx context.Context, bs blockstore.Blockstore, oldroot cid.Cid, newcids map[cid.Cid]blockformat.Block, skipcids map[cid.Cid]bool) (map[cid.Cid]bool, error) { ctx, span := otel.Tracer("repo").Start(ctx, "BlockDiff") defer span.End() @@ -1029,8 +805,8 @@ type UserStat struct { } func (cs *CarStore) Stat(ctx context.Context, usr models.Uid) ([]UserStat, error) { - var shards []CarShard - if err := cs.meta.Order("seq asc").Find(&shards, "usr = ?", usr).Error; err != nil { + shards, err := cs.meta.GetUserShards(ctx, usr) + if err != nil { return nil, err } @@ -1047,8 +823,8 @@ func (cs *CarStore) Stat(ctx context.Context, usr models.Uid) ([]UserStat, error } func (cs *CarStore) WipeUserData(ctx context.Context, user models.Uid) error { - var shards []*CarShard - if err := cs.meta.Find(&shards, "usr = ?", user).Error; err != nil { + shards, err := cs.meta.GetUserShards(ctx, user) + if err != nil { return err } @@ -1063,32 +839,23 @@ func (cs *CarStore) WipeUserData(ctx context.Context, user models.Uid) error { return nil } -func (cs *CarStore) deleteShards(ctx context.Context, shs []*CarShard) error { +func (cs *CarStore) deleteShards(ctx context.Context, shs []CarShard) error { ctx, span := otel.Tracer("carstore").Start(ctx, "deleteShards") defer span.End() - deleteSlice := func(ctx context.Context, subs []*CarShard) error { - var ids []uint - for _, sh := range subs { - ids = append(ids, sh.ID) - } - - txn := cs.meta.Begin() - - if err := txn.Delete(&CarShard{}, "id in (?)", ids).Error; err != nil { - return err - } - - if err := txn.Delete(&blockRef{}, "shard in (?)", ids).Error; err != nil { - return err + deleteSlice := func(ctx context.Context, subs []CarShard) error { + ids := make([]uint, len(subs)) + for i, sh := range subs { + ids[i] = sh.ID } - if err := txn.Commit().Error; err != nil { + err := cs.meta.DeleteShardsAndRefs(ctx, ids) + if err != nil { return err } for _, sh := range subs { - if err := cs.deleteShardFile(ctx, sh); err != nil { + if err := cs.deleteShardFile(ctx, &sh); err != nil { if !os.IsNotExist(err) { return err } @@ -1200,6 +967,8 @@ func (cb *compBucket) isEmpty() bool { func (cs *CarStore) openNewCompactedShardFile(ctx context.Context, user models.Uid, seq int) (*os.File, string, error) { // TODO: some overwrite protections + // NOTE CreateTemp is used for creating a non-colliding file, but we keep it and don't delete it so don't think of it as "temporary". + // This creates "sh-%d-%d%s" with some random stuff in the last position fi, err := os.CreateTemp(cs.rootDir, fnameForShard(user, seq)) if err != nil { return nil, "", err @@ -1217,31 +986,19 @@ func (cs *CarStore) GetCompactionTargets(ctx context.Context, shardCount int) ([ ctx, span := otel.Tracer("carstore").Start(ctx, "GetCompactionTargets") defer span.End() - var targets []CompactionTarget - if err := cs.meta.Raw(`select usr, count(*) as num_shards from car_shards group by usr having count(*) > ? order by num_shards desc`, shardCount).Scan(&targets).Error; err != nil { - return nil, err - } - - return targets, nil + return cs.meta.GetCompactionTargets(ctx, shardCount) } +// getBlockRefsForShards is a prep function for CompactUserShards func (cs *CarStore) getBlockRefsForShards(ctx context.Context, shardIds []uint) ([]blockRef, error) { ctx, span := otel.Tracer("carstore").Start(ctx, "getBlockRefsForShards") defer span.End() span.SetAttributes(attribute.Int("shards", len(shardIds))) - chunkSize := 2000 - out := make([]blockRef, 0, len(shardIds)) - for i := 0; i < len(shardIds); i += chunkSize { - sl := shardIds[i:] - if len(sl) > chunkSize { - sl = sl[:chunkSize] - } - - if err := blockRefsForShards(ctx, cs.meta, sl, &out); err != nil { - return nil, fmt.Errorf("getting block refs: %w", err) - } + out, err := cs.meta.GetBlockRefsForShards(ctx, shardIds) + if err != nil { + return nil, err } span.SetAttributes(attribute.Int("refs", len(out))) @@ -1249,31 +1006,6 @@ func (cs *CarStore) getBlockRefsForShards(ctx context.Context, shardIds []uint) return out, nil } -func valuesStatementForShards(shards []uint) string { - sb := new(strings.Builder) - for i, v := range shards { - sb.WriteByte('(') - sb.WriteString(strconv.Itoa(int(v))) - sb.WriteByte(')') - if i != len(shards)-1 { - sb.WriteByte(',') - } - } - return sb.String() -} - -func blockRefsForShards(ctx context.Context, db *gorm.DB, shards []uint, obuf *[]blockRef) error { - // Check the database driver - switch db.Dialector.(type) { - case *postgres.Dialector: - sval := valuesStatementForShards(shards) - q := fmt.Sprintf(`SELECT block_refs.* FROM block_refs INNER JOIN (VALUES %s) AS vals(v) ON block_refs.shard = v`, sval) - return db.Raw(q).Scan(obuf).Error - default: - return db.Raw(`SELECT * FROM block_refs WHERE shard IN (?)`, shards).Scan(obuf).Error - } -} - func shardSize(sh *CarShard) (int64, error) { st, err := os.Stat(sh.Path) if err != nil { @@ -1302,15 +1034,11 @@ func (cs *CarStore) CompactUserShards(ctx context.Context, user models.Uid, skip span.SetAttributes(attribute.Int64("user", int64(user))) - var shards []CarShard - if err := cs.meta.WithContext(ctx).Find(&shards, "usr = ?", user).Error; err != nil { + shards, err := cs.meta.GetUserShards(ctx, user) + if err != nil { return nil, err } - sort.Slice(shards, func(i, j int) bool { - return shards[i].Seq < shards[j].Seq - }) - if skipBigShards { // Since we generally expect shards to start bigger and get smaller, // and because we want to avoid compacting non-adjacent shards @@ -1356,8 +1084,8 @@ func (cs *CarStore) CompactUserShards(ctx context.Context, user models.Uid, skip span.SetAttributes(attribute.Int("blockRefs", len(brefs))) - var staleRefs []staleRef - if err := cs.meta.WithContext(ctx).Find(&staleRefs, "usr = ?", user).Error; err != nil { + staleRefs, err := cs.meta.GetUserStaleRefs(ctx, user) + if err != nil { return nil, err } @@ -1486,7 +1214,7 @@ func (cs *CarStore) CompactUserShards(ctx context.Context, user models.Uid, skip stats.NewShards++ - var todelete []*CarShard + todelete := make([]CarShard, 0, len(b.shards)) for _, s := range b.shards { removedShards[s.ID] = true sh, ok := shardsById[s.ID] @@ -1494,7 +1222,7 @@ func (cs *CarStore) CompactUserShards(ctx context.Context, user models.Uid, skip return nil, fmt.Errorf("missing shard to delete") } - todelete = append(todelete, &sh) + todelete = append(todelete, sh) } stats.ShardsDeleted += len(todelete) @@ -1546,27 +1274,7 @@ func (cs *CarStore) deleteStaleRefs(ctx context.Context, uid models.Uid, brefs [ } } - txn := cs.meta.Begin() - - if err := txn.Delete(&staleRef{}, "usr = ?", uid).Error; err != nil { - return err - } - - // now create a new staleRef with all the refs we couldn't clear out - if len(staleToKeep) > 0 { - if err := txn.Create(&staleRef{ - Usr: uid, - Cids: packCids(staleToKeep), - }).Error; err != nil { - return err - } - } - - if err := txn.Commit().Error; err != nil { - return fmt.Errorf("failed to commit staleRef updates: %w", err) - } - - return nil + return cs.meta.SetStaleRef(ctx, uid, staleToKeep) } func (cs *CarStore) compactBucket(ctx context.Context, user models.Uid, b *compBucket, shardsById map[uint]CarShard, keep map[cid.Cid]bool) error { diff --git a/carstore/meta_gorm.go b/carstore/meta_gorm.go new file mode 100644 index 00000000..eb9ff7bb --- /dev/null +++ b/carstore/meta_gorm.go @@ -0,0 +1,346 @@ +package carstore + +import ( + "bytes" + "context" + "fmt" + "io" + "strconv" + "strings" + "time" + + "github.com/bluesky-social/indigo/models" + "github.com/ipfs/go-cid" + "go.opentelemetry.io/otel" + "gorm.io/driver/postgres" + "gorm.io/gorm" +) + +type CarStoreGormMeta struct { + meta *gorm.DB +} + +func (cs *CarStoreGormMeta) Init() error { + if err := cs.meta.AutoMigrate(&CarShard{}, &blockRef{}); err != nil { + return err + } + if err := cs.meta.AutoMigrate(&staleRef{}); err != nil { + return err + } + return nil +} + +// Return true if any known record matches (Uid, Cid) +func (cs *CarStoreGormMeta) HasUidCid(ctx context.Context, user models.Uid, k cid.Cid) (bool, error) { + var count int64 + if err := cs.meta. + Model(blockRef{}). + Select("path, block_refs.offset"). + Joins("left join car_shards on block_refs.shard = car_shards.id"). + Where("usr = ? AND cid = ?", user, models.DbCID{CID: k}). + Count(&count).Error; err != nil { + return false, err + } + + return count > 0, nil +} + +// For some Cid, lookup the block ref. +// Return the path of the file written, the offset within the file, and the user associated with the Cid. +func (cs *CarStoreGormMeta) LookupBlockRef(ctx context.Context, k cid.Cid) (path string, offset int64, user models.Uid, err error) { + // TODO: for now, im using a join to ensure we only query blocks from the + // correct user. maybe it makes sense to put the user in the blockRef + // directly? tradeoff of time vs space + var info struct { + Path string + Offset int64 + Usr models.Uid + } + if err := cs.meta.Raw(`SELECT + (select path from car_shards where id = block_refs.shard) as path, + block_refs.offset, + (select usr from car_shards where id = block_refs.shard) as usr +FROM block_refs +WHERE + block_refs.cid = ? +LIMIT 1;`, models.DbCID{CID: k}).Scan(&info).Error; err != nil { + var defaultUser models.Uid + return "", -1, defaultUser, err + } + return info.Path, info.Offset, info.Usr, nil +} + +func (cs *CarStoreGormMeta) GetLastShard(ctx context.Context, user models.Uid) (*CarShard, error) { + var lastShard CarShard + if err := cs.meta.WithContext(ctx).Model(CarShard{}).Limit(1).Order("seq desc").Find(&lastShard, "usr = ?", user).Error; err != nil { + return nil, err + } + return &lastShard, nil +} + +// return all of a users's shards, ascending by Seq +func (cs *CarStoreGormMeta) GetUserShards(ctx context.Context, usr models.Uid) ([]CarShard, error) { + var shards []CarShard + if err := cs.meta.Order("seq asc").Find(&shards, "usr = ?", usr).Error; err != nil { + return nil, err + } + return shards, nil +} + +// return all of a users's shards, descending by Seq +func (cs *CarStoreGormMeta) GetUserShardsDesc(ctx context.Context, usr models.Uid, minSeq int) ([]CarShard, error) { + var shards []CarShard + if err := cs.meta.Order("seq desc").Find(&shards, "usr = ? AND seq >= ?", usr, minSeq).Error; err != nil { + return nil, err + } + return shards, nil +} + +func (cs *CarStoreGormMeta) GetUserStaleRefs(ctx context.Context, user models.Uid) ([]staleRef, error) { + var staleRefs []staleRef + if err := cs.meta.WithContext(ctx).Find(&staleRefs, "usr = ?", user).Error; err != nil { + return nil, err + } + return staleRefs, nil +} + +func (cs *CarStoreGormMeta) SeqForRev(ctx context.Context, user models.Uid, sinceRev string) (int, error) { + var untilShard CarShard + if err := cs.meta.Where("rev >= ? AND usr = ?", sinceRev, user).Order("rev").First(&untilShard).Error; err != nil { + return 0, fmt.Errorf("finding early shard: %w", err) + } + return untilShard.Seq, nil +} + +func (cs *CarStoreGormMeta) GetCompactionTargets(ctx context.Context, minShardCount int) ([]CompactionTarget, error) { + var targets []CompactionTarget + if err := cs.meta.Raw(`select usr, count(*) as num_shards from car_shards group by usr having count(*) > ? order by num_shards desc`, minShardCount).Scan(&targets).Error; err != nil { + return nil, err + } + + return targets, nil +} + +func (cs *CarStoreGormMeta) PutShardAndRefs(ctx context.Context, shard *CarShard, brefs []map[string]any, rmcids map[cid.Cid]bool) error { + // TODO: there should be a way to create the shard and block_refs that + // reference it in the same query, would save a lot of time + tx := cs.meta.WithContext(ctx).Begin() + + if err := tx.WithContext(ctx).Create(shard).Error; err != nil { + return fmt.Errorf("failed to create shard in DB tx: %w", err) + } + + for _, ref := range brefs { + ref["shard"] = shard.ID + } + + if err := createBlockRefs(ctx, tx, brefs); err != nil { + return fmt.Errorf("failed to create block refs: %w", err) + } + + if len(rmcids) > 0 { + cids := make([]cid.Cid, 0, len(rmcids)) + for c := range rmcids { + cids = append(cids, c) + } + + if err := tx.Create(&staleRef{ + Cids: packCids(cids), + Usr: shard.Usr, + }).Error; err != nil { + return err + } + } + + err := tx.WithContext(ctx).Commit().Error + if err != nil { + return fmt.Errorf("failed to commit shard DB transaction: %w", err) + } + return nil +} + +func (cs *CarStoreGormMeta) DeleteShardsAndRefs(ctx context.Context, ids []uint) error { + txn := cs.meta.Begin() + + if err := txn.Delete(&CarShard{}, "id in (?)", ids).Error; err != nil { + txn.Rollback() + return err + } + + if err := txn.Delete(&blockRef{}, "shard in (?)", ids).Error; err != nil { + txn.Rollback() + return err + } + + return txn.Commit().Error +} + +func (cs *CarStoreGormMeta) GetBlockRefsForShards(ctx context.Context, shardIds []uint) ([]blockRef, error) { + chunkSize := 2000 + out := make([]blockRef, 0, len(shardIds)) + for i := 0; i < len(shardIds); i += chunkSize { + sl := shardIds[i:] + if len(sl) > chunkSize { + sl = sl[:chunkSize] + } + + if err := blockRefsForShards(ctx, cs.meta, sl, &out); err != nil { + return nil, fmt.Errorf("getting block refs: %w", err) + } + } + return out, nil +} + +// blockRefsForShards is an inner loop helper for GetBlockRefsForShards +func blockRefsForShards(ctx context.Context, db *gorm.DB, shards []uint, obuf *[]blockRef) error { + // Check the database driver + switch db.Dialector.(type) { + case *postgres.Dialector: + sval := valuesStatementForShards(shards) + q := fmt.Sprintf(`SELECT block_refs.* FROM block_refs INNER JOIN (VALUES %s) AS vals(v) ON block_refs.shard = v`, sval) + return db.Raw(q).Scan(obuf).Error + default: + return db.Raw(`SELECT * FROM block_refs WHERE shard IN (?)`, shards).Scan(obuf).Error + } +} + +// valuesStatementForShards builds a postgres compatible statement string from int literals +func valuesStatementForShards(shards []uint) string { + sb := new(strings.Builder) + for i, v := range shards { + sb.WriteByte('(') + sb.WriteString(strconv.Itoa(int(v))) + sb.WriteByte(')') + if i != len(shards)-1 { + sb.WriteByte(',') + } + } + return sb.String() +} + +func (cs *CarStoreGormMeta) SetStaleRef(ctx context.Context, uid models.Uid, staleToKeep []cid.Cid) error { + txn := cs.meta.Begin() + + if err := txn.Delete(&staleRef{}, "usr = ?", uid).Error; err != nil { + return err + } + + // now create a new staleRef with all the refs we couldn't clear out + if len(staleToKeep) > 0 { + if err := txn.Create(&staleRef{ + Usr: uid, + Cids: packCids(staleToKeep), + }).Error; err != nil { + return err + } + } + + if err := txn.Commit().Error; err != nil { + return fmt.Errorf("failed to commit staleRef updates: %w", err) + } + return nil +} + +type CarShard struct { + ID uint `gorm:"primarykey"` + CreatedAt time.Time + + Root models.DbCID `gorm:"index"` + DataStart int64 + Seq int `gorm:"index:idx_car_shards_seq;index:idx_car_shards_usr_seq,priority:2,sort:desc"` + Path string + Usr models.Uid `gorm:"index:idx_car_shards_usr;index:idx_car_shards_usr_seq,priority:1"` + Rev string +} + +type blockRef struct { + ID uint `gorm:"primarykey"` + Cid models.DbCID `gorm:"index"` + Shard uint `gorm:"index"` + Offset int64 + //User uint `gorm:"index"` +} + +type staleRef struct { + ID uint `gorm:"primarykey"` + Cid *models.DbCID + Cids []byte + Usr models.Uid `gorm:"index"` +} + +func (sr *staleRef) getCids() ([]cid.Cid, error) { + if sr.Cid != nil { + return []cid.Cid{sr.Cid.CID}, nil + } + + return unpackCids(sr.Cids) +} + +func unpackCids(b []byte) ([]cid.Cid, error) { + br := bytes.NewReader(b) + var out []cid.Cid + for { + _, c, err := cid.CidFromReader(br) + if err != nil { + if err == io.EOF { + break + } + return nil, err + } + + out = append(out, c) + } + + return out, nil +} + +func packCids(cids []cid.Cid) []byte { + buf := new(bytes.Buffer) + for _, c := range cids { + buf.Write(c.Bytes()) + } + + return buf.Bytes() +} + +func createBlockRefs(ctx context.Context, tx *gorm.DB, brefs []map[string]any) error { + ctx, span := otel.Tracer("carstore").Start(ctx, "createBlockRefs") + defer span.End() + + if err := createInBatches(ctx, tx, brefs, 2000); err != nil { + return err + } + + return nil +} + +// Function to create in batches +func createInBatches(ctx context.Context, tx *gorm.DB, brefs []map[string]any, batchSize int) error { + for i := 0; i < len(brefs); i += batchSize { + batch := brefs[i:] + if len(batch) > batchSize { + batch = batch[:batchSize] + } + + query, values := generateInsertQuery(batch) + + if err := tx.WithContext(ctx).Exec(query, values...).Error; err != nil { + return err + } + } + return nil +} + +func generateInsertQuery(brefs []map[string]any) (string, []any) { + placeholders := strings.Repeat("(?, ?, ?),", len(brefs)) + placeholders = placeholders[:len(placeholders)-1] // trim trailing comma + + query := "INSERT INTO block_refs (\"cid\", \"offset\", \"shard\") VALUES " + placeholders + + values := make([]any, 0, 3*len(brefs)) + for _, entry := range brefs { + values = append(values, entry["cid"], entry["offset"], entry["shard"]) + } + + return query, values +} diff --git a/carstore/util.go b/carstore/util.go new file mode 100644 index 00000000..56396501 --- /dev/null +++ b/carstore/util.go @@ -0,0 +1,32 @@ +package carstore + +import ( + "encoding/binary" + "io" +) + +// Length-delimited Write +// Writer stream gets Uvarint length then concatenated data +func LdWrite(w io.Writer, d ...[]byte) (int64, error) { + var sum uint64 + for _, s := range d { + sum += uint64(len(s)) + } + + buf := make([]byte, 8) + n := binary.PutUvarint(buf, sum) + nw, err := w.Write(buf[:n]) + if err != nil { + return 0, err + } + + for _, s := range d { + onw, err := w.Write(s) + if err != nil { + return int64(nw), err + } + nw += onw + } + + return int64(nw), nil +}