storj/storagenode/storagenodedb/usedserials.go
ethanadams 47e4584fbe
V3-1989: Storage node database is locked for several minutes while submiting orders (#2410)
* remove infodb locks and give a unique name for each in memory created.

* changed max idle and open to 1 for memory DBs.  fixes table locking errors

* fixed race condition

* added file based infodb test

* added busy timeout parameter to the file based infodb for testing

* fixed imports

* removed db.locked() after merge from master
2019-07-02 17:23:02 -04:00

72 lines
1.9 KiB
Go

// Copyright (C) 2019 Storj Labs, Inc.
// See LICENSE for copying information.
package storagenodedb
import (
"context"
"time"
"github.com/zeebo/errs"
"storj.io/storj/pkg/storj"
"storj.io/storj/storagenode/piecestore"
)
type usedSerials struct {
*InfoDB
}
// UsedSerials returns certificate database.
func (db *DB) UsedSerials() piecestore.UsedSerials { return db.info.UsedSerials() }
// UsedSerials returns certificate database.
func (db *InfoDB) UsedSerials() piecestore.UsedSerials { return &usedSerials{db} }
// Add adds a serial to the database.
func (db *usedSerials) Add(ctx context.Context, satelliteID storj.NodeID, serialNumber storj.SerialNumber, expiration time.Time) (err error) {
defer mon.Task()(&ctx)(&err)
_, err = db.db.Exec(`
INSERT INTO
used_serial(satellite_id, serial_number, expiration)
VALUES(?, ?, ?)`, satelliteID, serialNumber, expiration)
return ErrInfo.Wrap(err)
}
// DeleteExpired deletes expired serial numbers
func (db *usedSerials) DeleteExpired(ctx context.Context, now time.Time) (err error) {
defer mon.Task()(&ctx)(&err)
_, err = db.db.Exec(`DELETE FROM used_serial WHERE expiration < ?`, now)
return ErrInfo.Wrap(err)
}
// IterateAll iterates all serials.
// Note, this will lock the database and should only be used during startup.
func (db *usedSerials) IterateAll(ctx context.Context, fn piecestore.SerialNumberFn) (err error) {
defer mon.Task()(&ctx)(&err)
rows, err := db.db.Query(`SELECT satellite_id, serial_number, expiration FROM used_serial`)
if err != nil {
return ErrInfo.Wrap(err)
}
defer func() { err = errs.Combine(err, ErrInfo.Wrap(rows.Close())) }()
for rows.Next() {
var satelliteID storj.NodeID
var serialNumber storj.SerialNumber
var expiration time.Time
err := rows.Scan(&satelliteID, &serialNumber, &expiration)
if err != nil {
return ErrInfo.Wrap(err)
}
fn(satelliteID, serialNumber, expiration)
}
return ErrInfo.Wrap(rows.Err())
}