buildbot/nix-eval Build done.
PR Size Review Check / pr-size-review-gate (pull_request) Failing after 51s
buildbot/nix-build gitea:clan/data-mesher#checks.x86_64-linux.nixos-data-mesher-basic Build done.
sizelint / sizelint (pull_request) Successful in 1m8s
buildbot/nix-build Build done.
Makes it simpler to integrate with clan service module and to reason about which files belong to which network by examining the file system.
273 lines
6.3 KiB
Go
273 lines
6.3 KiB
Go
package state
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"path"
|
|
|
|
"git.clan.lol/clan/data-mesher/pkg/config"
|
|
"git.clan.lol/clan/data-mesher/pkg/crypto"
|
|
"git.clan.lol/clan/data-mesher/pkg/model"
|
|
"git.clan.lol/clan/data-mesher/pkg/msgpack"
|
|
bolt "go.etcd.io/bbolt"
|
|
)
|
|
|
|
var (
|
|
bucketSignatures = []byte("signatures") //nolint:gochecknoglobals
|
|
|
|
ErrNetworkNotFound = errors.New("network not found")
|
|
ErrSignatureNotFound = errors.New("signature not found")
|
|
ErrNetworkNotConfigured = errors.New("network not configured")
|
|
)
|
|
|
|
type SignatureReader interface {
|
|
Read(sigs []model.Signature) (int, error)
|
|
}
|
|
|
|
type signatureReader struct {
|
|
tx *bolt.Tx
|
|
networks []*crypto.PublicKey
|
|
|
|
cursor *bolt.Cursor
|
|
cursorStarted bool
|
|
networkIdx int
|
|
}
|
|
|
|
func newSignatureReader(tx *bolt.Tx) (*signatureReader, error) {
|
|
sigsBucket := tx.Bucket(bucketSignatures)
|
|
|
|
// get a list of all networks
|
|
var networks []*crypto.PublicKey
|
|
|
|
err := sigsBucket.ForEachBucket(func(name []byte) error {
|
|
id, err := crypto.ParsePublicKey(name)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to parse network ID: %w", err)
|
|
}
|
|
|
|
networks = append(networks, id)
|
|
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to list signature buckets: %w", err)
|
|
}
|
|
|
|
return &signatureReader{
|
|
tx: tx,
|
|
networks: networks,
|
|
}, nil
|
|
}
|
|
|
|
func (r *signatureReader) Read(sigs []model.Signature) (int, error) {
|
|
var n int
|
|
|
|
for n < len(sigs) {
|
|
// advance to the next network bucket if needed
|
|
if r.cursor == nil {
|
|
if r.networkIdx >= len(r.networks) {
|
|
return n, io.EOF
|
|
}
|
|
|
|
bucket := r.tx.Bucket(bucketSignatures).Bucket(r.networks[r.networkIdx].Bytes())
|
|
r.cursor = bucket.Cursor()
|
|
r.cursorStarted = false
|
|
r.networkIdx++
|
|
}
|
|
|
|
var k, v []byte
|
|
if !r.cursorStarted {
|
|
k, v = r.cursor.First()
|
|
r.cursorStarted = true
|
|
} else {
|
|
k, v = r.cursor.Next()
|
|
}
|
|
|
|
if k == nil {
|
|
// end of this bucket, try next network on next iteration
|
|
r.cursor = nil
|
|
|
|
continue
|
|
}
|
|
|
|
if err := msgpack.Unmarshal(v, &sigs[n]); err != nil {
|
|
return n, fmt.Errorf("failed to unmarshal signature: %w", err)
|
|
}
|
|
|
|
n++
|
|
}
|
|
|
|
return n, nil
|
|
}
|
|
|
|
type Signatures struct {
|
|
db *bolt.DB
|
|
}
|
|
|
|
func NewSignatures(cfg *config.Config) (*Signatures, error) {
|
|
// open the database
|
|
db, err := bolt.Open(path.Join(cfg.StateDirectory, "signatures.db"), 0o600, nil)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to open file metadata database: %w", err)
|
|
}
|
|
|
|
// ensure the top-level signatures bucket exists
|
|
err = db.Update(func(tx *bolt.Tx) error {
|
|
_, createErr := tx.CreateBucketIfNotExists(bucketSignatures)
|
|
if createErr != nil {
|
|
return fmt.Errorf("failed to create bucket: %w", createErr)
|
|
}
|
|
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to create signatures bucket: %w", err)
|
|
}
|
|
|
|
return &Signatures{db}, nil
|
|
}
|
|
|
|
func (s *Signatures) BeginTx(writable bool) (*bolt.Tx, error) {
|
|
tx, err := s.db.Begin(writable)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to begin transaction: %w", err)
|
|
}
|
|
|
|
return tx, nil
|
|
}
|
|
|
|
func (s *Signatures) Networks(tx *bolt.Tx) ([]*crypto.PublicKey, error) {
|
|
var result []*crypto.PublicKey
|
|
|
|
sigsBucket := tx.Bucket(bucketSignatures)
|
|
|
|
// iterate each child bucket, 1 per network
|
|
err := sigsBucket.ForEachBucket(func(name []byte) error {
|
|
id, parseErr := crypto.ParsePublicKey(name)
|
|
if parseErr != nil {
|
|
return fmt.Errorf("failed to parse network ID: %w", parseErr)
|
|
}
|
|
|
|
result = append(result, id)
|
|
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to iterate over buckets: %w", err)
|
|
}
|
|
|
|
return result, nil
|
|
}
|
|
|
|
func (s *Signatures) List(tx *bolt.Tx) (SignatureReader, error) {
|
|
return newSignatureReader(tx)
|
|
}
|
|
|
|
func (s *Signatures) Get(tx *bolt.Tx, network *crypto.PublicKey, name string, sig *model.Signature) error {
|
|
bucket, err := s.bucketForNetwork(tx, network)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
buf := bucket.Get([]byte(name))
|
|
if buf == nil {
|
|
return fmt.Errorf("%w: %s", ErrSignatureNotFound, sig.Name)
|
|
}
|
|
|
|
if err := msgpack.Unmarshal(buf, sig); err != nil {
|
|
return fmt.Errorf("failed to unmarshal signature: %w", err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Signatures) Put(tx *bolt.Tx, sig *model.Signature) error {
|
|
// we partition signatures by network ID
|
|
bucket, err := s.ensureBucketForNetwork(tx, sig.NetworkID)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to create bucket for network %s: %w", sig.NetworkID, err)
|
|
}
|
|
|
|
buf, err := msgpack.Marshal(sig)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to marshal signature: %w", err)
|
|
}
|
|
|
|
if err = bucket.Put([]byte(sig.Name), buf); err != nil {
|
|
return fmt.Errorf("failed to put signature for network %s, name %s: %w", sig.NetworkID, sig.Name, err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Signatures) PutIfLater(tx *bolt.Tx, sig *model.Signature) (bool, error) {
|
|
// ensure the bucket exists
|
|
if _, err := s.ensureBucketForNetwork(tx, sig.NetworkID); err != nil {
|
|
return false, fmt.Errorf("failed to ensure bucket for network %s: %w", sig.NetworkID, err)
|
|
}
|
|
|
|
var (
|
|
updated bool
|
|
current model.Signature
|
|
)
|
|
|
|
err := s.Get(tx, sig.NetworkID, sig.Name, ¤t)
|
|
|
|
switch {
|
|
case errors.Is(err, ErrSignatureNotFound):
|
|
err = s.Put(tx, sig)
|
|
updated = true
|
|
case err != nil:
|
|
return false, fmt.Errorf("failed to get signature: %w", err)
|
|
case sig.SignedAt.After(current.SignedAt):
|
|
err = s.Put(tx, sig)
|
|
updated = true
|
|
}
|
|
|
|
return updated, err
|
|
}
|
|
|
|
func (s *Signatures) Delete(tx *bolt.Tx, network *crypto.PublicKey, name string) error {
|
|
bucket, err := s.bucketForNetwork(tx, network)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if err = bucket.Delete([]byte(name)); err != nil {
|
|
return fmt.Errorf("failed to delete key %s: %w", name, err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Signatures) Close() error {
|
|
if err := s.db.Close(); err != nil {
|
|
return fmt.Errorf("failed to close database: %w", err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Signatures) bucketForNetwork(tx *bolt.Tx, network *crypto.PublicKey) (*bolt.Bucket, error) {
|
|
sigsBucket := tx.Bucket(bucketSignatures)
|
|
|
|
bucket := sigsBucket.Bucket(network.Bytes())
|
|
if bucket == nil {
|
|
return nil, fmt.Errorf("%w: %s", ErrNetworkNotFound, network)
|
|
}
|
|
|
|
return bucket, nil
|
|
}
|
|
|
|
func (s *Signatures) ensureBucketForNetwork(tx *bolt.Tx, network *crypto.PublicKey) (*bolt.Bucket, error) {
|
|
sigsBucket := tx.Bucket(bucketSignatures)
|
|
|
|
bucket, err := sigsBucket.CreateBucketIfNotExists(network.Bytes())
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to create bucket for network %s: %w", network, err)
|
|
}
|
|
|
|
return bucket, nil
|
|
}
|