Files
brianmcgee 48c07d47d2
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.
state: partition files on disk by network name instead of network id
Makes it simpler to integrate with clan service module and to reason about which files belong to which network by examining the file system.
2026-03-02 11:50:18 +00:00

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, &current)
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
}