add TransactionManager and TransactionStore for transactions/backups

This all needs to be wired into API/Server/Cluster/Holder etc. but I
think the TransactionManager will be a pretty good building block for
managing transaction state at the coordinator level.
This commit is contained in:
Matt Jaffee 2020-04-15 17:04:14 -05:00
parent c3ef9a1768
commit 38cec6f20e
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
2 changed files with 623 additions and 0 deletions

359
transaction.go Normal file
View file

@ -0,0 +1,359 @@
package pilosa
import (
"sync"
"time"
"github.com/pilosa/pilosa/v2/logger"
"github.com/pkg/errors"
)
// Transaction contains information related to a block of work that
// needs to be tracked and spans multiple API calls.
type Transaction struct {
// ID is an arbitrary string identifier. All transactions must have a unique ID.
ID string
// Active notes whether an Exclusive transaction is active, or
// still pending (if other active transactions exist). All
// non-exclusive transactions are always active.
Active bool
// Exclusive is set on transactions which can only become active when no other transactions exist.
Exclusive bool
// Timeout is the minimum idle time for which this transaction should continue to exist.
Timeout time.Duration
// Deadline is calculated from Timeout, and should be reset each
// time there is activity on the transaction.
Deadline time.Time
// Stats track statistics for the transaction. Not yet used.
Stats TransactionStats
}
type TransactionStats struct{}
// TransactionManager enforces the rules for transactions on a single
// node. It is goroutine-safe. It should be created by a call to
// NewTransactionManager where it takes a TransactionStore. If logging
// is desired, Log should be set before an instance of
// TransactionManager is used.
type TransactionManager struct {
mu sync.RWMutex
Log logger.Logger
store TransactionStore
checkingDeadlines bool
}
// NewTransactionManager creates a new TransactionManager with the
// given store.
func NewTransactionManager(store TransactionStore) *TransactionManager {
tm := &TransactionManager{
Log: logger.NopLogger,
store: store,
checkingDeadlines: true,
}
// start deadline checker in case we've just started up, but there is already state in the store.
go tm.deadlineChecker()
return tm
}
// Start starts a new transaction with the given parameters. If an
// exclusive transaction is pending or in progress,
// ErrTransactionExclusive is returned. If a transaction with the same
// id already exists, that transaction is returned along with
// ErrTransactionExists. If there is no error, the created transaction
// is returned—this is primarily so that the caller can discover if an
// exclusive transaction has been made immediately active or if they
// need to poll.
func (tm *TransactionManager) Start(id string, timeout time.Duration, exclusive bool) (Transaction, error) {
tm.mu.Lock()
defer tm.mu.Unlock()
trnsMap, err := tm.store.List()
if err != nil {
return Transaction{}, errors.Wrap(err, "listing transactions in Start")
}
// check for an exclusive transaction
for _, trns := range trnsMap {
if trns.Exclusive {
// if someone wants a transaction, and we're not able to
// give it to them, we want to be checking deadlines.
tm.startDeadlineChecker()
return Transaction{}, ErrTransactionExclusive
}
}
if trns, ok := trnsMap[id]; ok {
return trns, ErrTransactionExists
}
// set new transaction to active if it is not exclusive or if
// there are no other transactions.
active := !exclusive || (len(trnsMap) == 0)
// set deadline according to timeout
deadline := time.Now().Add(timeout)
trns := Transaction{
ID: id,
Active: active,
Exclusive: exclusive,
Timeout: timeout,
Deadline: deadline,
}
err = tm.store.Put(trns)
// we won't check deadlines unless there's actually an exclusive
// transaction pending
if exclusive && !active {
tm.startDeadlineChecker()
}
return trns, errors.Wrap(err, "adding to store")
}
// Finish completes and removes a transaction, returning the completed
// transaction (so that the caller can e.g. view the Stats)
func (tm *TransactionManager) Finish(id string) (Transaction, error) {
tm.mu.Lock()
defer tm.mu.Unlock()
return tm.finish(id)
}
// finish is the unprotected implementation of Finish
func (tm *TransactionManager) finish(id string) (Transaction, error) {
// sanity check
if trns, err := tm.store.Get(id); err != nil {
return trns, err
}
trns, err := tm.store.Remove(id)
if err != nil {
return trns, err
}
// After removing, check to see if we need to activate an exclusive transaction
trnsMap, err := tm.store.List()
if err != nil {
// returning an error here is weird because we've already
// removed the transaction
return trns, errors.Wrap(err, "listing transactions in Finish")
}
if len(trnsMap) == 1 {
for _, etrans := range trnsMap {
if etrans.Exclusive {
if etrans.Active { // sanity check
panic("we just removed a transaction, and the sole remaining exclusive transaction was already active")
}
etrans.Active = true
etrans.Deadline = time.Now().Add(etrans.Timeout)
if err := tm.store.Put(etrans); err != nil {
return trns, errors.Wrap(err, "activating exclusive transaction after finishing last transaction")
}
}
}
}
return trns, nil
}
// Get retrieves the transaction with the given ID. Returns ErrTransactionNotFound
// if there isn't one.
func (tm *TransactionManager) Get(id string) (Transaction, error) {
tm.mu.RLock()
tm.mu.RUnlock()
return tm.store.Get(id)
}
// List returns map of all transactions by their ID. It is a copy and
// so may be retained and modified by the caller.
func (tm *TransactionManager) List() (map[string]Transaction, error) {
tm.mu.RLock()
defer tm.mu.RUnlock()
return tm.store.List()
}
// ResetDeadline updates the deadline for the transaction with the
// given ID to be equal to the current time plus the transaction's
// timeout.
func (tm *TransactionManager) ResetDeadline(id string) (Transaction, error) {
tm.mu.Lock()
defer tm.mu.Unlock()
trns, err := tm.store.Get(id)
if err != nil {
return trns, errors.Wrap(err, "getting transaction")
}
trns.Deadline = time.Now().Add(trns.Timeout)
err = tm.store.Put(trns)
return trns, errors.Wrap(err, "storing transaction with new timeout")
}
// startDeadlineChecker may only be called while tm.mu is held.
func (tm *TransactionManager) startDeadlineChecker() {
if !tm.checkingDeadlines {
tm.checkingDeadlines = true
go tm.deadlineChecker()
}
}
// deadlineChecker loops continuously checking for expired
// deadlines. It stops when there are no upcoming deadlines.
func (tm *TransactionManager) deadlineChecker() {
interval := tm.checkDeadlines()
for interval != 0 {
time.Sleep(interval)
interval = tm.checkDeadlines()
}
tm.mu.Lock()
tm.checkingDeadlines = false
tm.mu.Unlock()
}
// checkDeadlines finishes transactions which are past their
// deadlines. It returns the duration until the next deadline. If
// there are no exclusive transactions, it does nothing and returns 0
// as a signal to stop checking.
func (tm *TransactionManager) checkDeadlines() time.Duration {
tm.mu.Lock()
defer tm.mu.Unlock()
trnsMap, err := tm.store.List()
if err != nil {
tm.log().Printf("transaction deadline checker couldn't list transactions: %v", err)
return 0
}
hasExclusive := false
for _, trns := range trnsMap {
if trns.Exclusive {
hasExclusive = true
break
}
}
if !hasExclusive {
return 0 // no need to expire things if nothing is waiting
}
now := time.Now()
// track the time interval to next deadline
nextInterval := time.Duration(0)
for id, trns := range trnsMap {
// fmt.Printf("trns: %v", id)
if !trns.Active {
// fmt.Printf(" not active\n")
continue
}
if !now.Before(trns.Deadline) {
// fmt.Printf(" finishing\n")
trnsF, err := tm.finish(id)
if err != nil {
tm.log().Printf("error finishing expired transaction '%s': %+v: %v", id, trnsF, err)
} else {
tm.log().Printf("cleared expired transaction: %+v", trnsF)
}
} else {
interval := trns.Deadline.Sub(now)
// fmt.Printf(" getting new interval: %v, next: %v\n", interval, nextInterval)
if nextInterval == 0 || interval < nextInterval {
nextInterval = interval
}
}
}
return nextInterval
}
func (tm *TransactionManager) log() logger.Logger {
if tm.Log != nil {
return tm.Log
}
return logger.NopLogger
}
// TransactionStore declares the functionality which a store for
// Pilosa transactions must implement.
type TransactionStore interface {
// Put stores a new transaction or replaces an existing transaction with the given one.
Put(trns Transaction) error
// Get retrieves the transaction at id or returns ErrTransactionNotFound if there isn't one.
Get(id string) (Transaction, error)
// List returns a map of all transactions by ID. The map must be safe to modify by the caller.
List() (map[string]Transaction, error)
// Remove deletes the transaction from the store. It must return ErrTransactionNotFound if there isn't one.
Remove(id string) (Transaction, error)
}
type OpenTransactionStoreFunc func(path string) (TransactionStore, error)
func OpenInMemTransactionStore(path string) (TransactionStore, error) {
return NewInMemTransactionStore(), nil
}
// InMemTransactionStore does not persist transaction data and is only
// useful for testing.
type InMemTransactionStore struct {
mu sync.RWMutex
tmap map[string]Transaction
}
func NewInMemTransactionStore() *InMemTransactionStore {
return &InMemTransactionStore{
tmap: make(map[string]Transaction),
}
}
func (s *InMemTransactionStore) Put(trns Transaction) error {
s.mu.Lock()
defer s.mu.Unlock()
s.tmap[trns.ID] = trns
return nil
}
func (s *InMemTransactionStore) Get(id string) (Transaction, error) {
s.mu.RLock()
defer s.mu.RUnlock()
if trns, ok := s.tmap[id]; ok {
return trns, nil
} else {
return Transaction{}, ErrTransactionNotFound
}
}
func (s *InMemTransactionStore) List() (map[string]Transaction, error) {
cp := make(map[string]Transaction)
for id, trns := range s.tmap {
cp[id] = trns
}
return cp, nil
}
func (s *InMemTransactionStore) Remove(id string) (Transaction, error) {
s.mu.Lock()
defer s.mu.Unlock()
if trns, ok := s.tmap[id]; ok {
delete(s.tmap, id)
return trns, nil
} else {
return Transaction{}, ErrTransactionNotFound
}
}
type Error string
func (e Error) Error() string { return string(e) }
const ErrTransactionNotFound = Error("transaction not found")
const ErrTransactionExclusive = Error("there is already an exclusive transaction")
const ErrTransactionExists = Error("transaction with the given id already exists")
const ErrTransactionInactive = Error("cannot finish an inactive transaction")

264
transaction_test.go Normal file
View file

@ -0,0 +1,264 @@
package pilosa_test
import (
"testing"
"time"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/test"
)
// TestTransactionManager currently uses an in memory transaction
// store, but tests a variety of timeouts, and therefore could be
// sensitive to slowness in the implementation. Especially if a store
// were used that actually wrote things to disk.
func TestTransactionManager(t *testing.T) {
store := pilosa.NewInMemTransactionStore()
tm := pilosa.NewTransactionManager(store)
tm.Log = test.NewBufferLogger()
// can add a non-exclusive transaction
trns1 := mustStart(t, tm, "a", time.Microsecond, false)
compareTransactions(t, pilosa.Transaction{ID: "a", Active: true, Timeout: time.Microsecond, Deadline: time.Now()}, trns1)
// can have two non exclusive transactions
trns2 := mustStart(t, tm, "b", time.Microsecond, false)
compareTransactions(t, pilosa.Transaction{ID: "b", Active: true, Timeout: time.Microsecond, Deadline: time.Now()}, trns2)
// trying to start a transaction with same name errors and returns previous transaction
t3, err := tm.Start("a", time.Second, true)
if err != pilosa.ErrTransactionExists {
t.Errorf("expected transaction exists, but got: '%v'", err)
}
compareTransactions(t, trns1, t3)
// can get an existing transaction
trns2_2 := mustGet(t, tm, "b")
compareTransactions(t, trns2, trns2_2)
// can list all transactions
trnsMap := mustList(t, tm)
if len(trnsMap) != 2 {
t.Errorf("unexpected number of transactions in map: %d", len(trnsMap))
}
compareTransactions(t, trnsMap["a"], trns1)
compareTransactions(t, trnsMap["b"], trns2)
// can submit an exclusive transaction
trnsE := mustStart(t, tm, "ce", time.Millisecond*5, true)
compareTransactions(t, pilosa.Transaction{ID: "ce", Active: false, Exclusive: true, Timeout: time.Millisecond * 5, Deadline: time.Now().Add(time.Millisecond * 5)}, trnsE)
// can't start new transactions while an exclusive transaction is pending
if _, err := tm.Start("d", time.Millisecond, false); err != pilosa.ErrTransactionExclusive {
t.Errorf("unexpected error starting transaction while an exclusive transaction exists: %v", err)
}
// can't start new exclusive transactions while an exclusive transaction is pending
if _, err := tm.Start("ee", time.Millisecond, true); err != pilosa.ErrTransactionExclusive {
t.Errorf("unexpected error starting transaction while an exclusive transaction exists: %v", err)
}
// exclusive transaction becomes active after deadlines expire
for i := 0; true; i++ {
time.Sleep(time.Microsecond)
trnsE, err := tm.Get("ce")
if err != nil {
t.Errorf("error retrieving exclusive transaction: %v", err)
}
if trnsE.Active {
break
}
if i > 100 {
t.Fatalf("exclusive transaction never became active: %+v", trnsE)
}
}
// can't start new transactions while an exclusive transaction is active
if _, err := tm.Start("f", time.Millisecond, false); err != pilosa.ErrTransactionExclusive {
t.Errorf("unexpected error starting transaction while an exclusive transaction exists: %v", err)
}
// can't start new exclusive transactions while an exclusive transaction is active
if _, err := tm.Start("ge", time.Millisecond, true); err != pilosa.ErrTransactionExclusive {
t.Errorf("unexpected error starting transaction while an exclusive transaction exists: %v", err)
}
// exclusive transaction gets expired after other transactions have attempted to start
for i := 0; true; i++ {
time.Sleep(time.Millisecond * 2)
trnsE, err := tm.Get("ce")
if err == nil {
if i > 10 {
t.Fatalf("exclusive transaction didn't expire: %+v", trnsE)
}
} else if err != pilosa.ErrTransactionNotFound {
t.Errorf("unexpected error fetching transaction while waiting for expiration: %v", err)
} else {
break // transaction was not found, therefore it expired and we can happily continue
}
}
// can start a new exclusive transaction and it's immediately active
trnsHE := mustStart(t, tm, "he", time.Hour, true)
compareTransactions(t, pilosa.Transaction{ID: "he", Active: true, Exclusive: true, Timeout: time.Hour, Deadline: time.Now().Add(time.Hour)}, trnsHE)
// can't start new transactions while an exclusive transaction is active
if _, err := tm.Start("i", time.Millisecond, false); err != pilosa.ErrTransactionExclusive {
t.Errorf("unexpected error starting transaction while an exclusive transaction exists: %v", err)
}
// can finish an active exclusive transaction
trnsHE_finish := mustFinish(t, tm, "he")
compareTransactions(t, trnsHE, trnsHE_finish)
// can start normal transaction after finishing exclusive transaction
trnsJ := mustStart(t, tm, "j", time.Hour, false)
compareTransactions(t, pilosa.Transaction{ID: "j", Active: true, Timeout: time.Hour, Deadline: time.Now().Add(time.Hour)}, trnsJ)
// can finish normal transaction
trnsJ_finish := mustFinish(t, tm, "j")
compareTransactions(t, trnsJ, trnsJ_finish)
// can start normal transaction after finishing normal transaction
trnsK := mustStart(t, tm, "k", time.Hour, false)
compareTransactions(t, pilosa.Transaction{ID: "k", Active: true, Timeout: time.Hour, Deadline: time.Now().Add(time.Hour)}, trnsK)
// can start new exclusive transaction, but not immediately active
trnsLE := mustStart(t, tm, "le", time.Hour, true)
compareTransactions(t, pilosa.Transaction{ID: "le", Exclusive: true, Timeout: time.Hour, Deadline: time.Now().Add(time.Hour)}, trnsLE)
// finishing k should activate le
trnsK_finish := mustFinish(t, tm, "k")
compareTransactions(t, trnsK, trnsK_finish)
trnsLE_active := mustGet(t, tm, "le")
trnsLE.Active = true
compareTransactions(t, trnsLE, trnsLE_active)
mustFinish(t, tm, "le")
// can start normal transaction to test deadline reset
trnsM := mustStart(t, tm, "m", time.Millisecond*4, false)
compareTransactions(t, pilosa.Transaction{ID: "m", Active: true, Timeout: time.Millisecond * 4, Deadline: time.Now().Add(time.Millisecond * 4)}, trnsM)
// start new exclusive transaction to trigger deadline check
trnsNE := mustStart(t, tm, "ne", time.Hour, true)
compareTransactions(t, pilosa.Transaction{ID: "ne", Exclusive: true, Timeout: time.Hour, Deadline: time.Now().Add(time.Hour)}, trnsNE)
// sleep for most of the deadline
time.Sleep(time.Millisecond * 3)
// reset deadline
trnsM_reset, err := tm.ResetDeadline("m")
if err != nil {
t.Errorf("resetting deadline: %v", err)
}
trnsM.Deadline = time.Now().Add(time.Millisecond * 4)
compareTransactions(t, trnsM, trnsM_reset)
// sleep until past the original deadline
time.Sleep(time.Millisecond * 2)
// verify that trnsM still exists
trnsM_again := mustGet(t, tm, "m")
compareTransactions(t, trnsM, trnsM_again)
}
func mustStart(t *testing.T, tm *pilosa.TransactionManager, id string, timeout time.Duration, exclusive bool) pilosa.Transaction {
t.Helper()
trns, err := tm.Start(id, timeout, exclusive)
if err != nil {
t.Errorf("starting transaction: %v", err)
}
return trns
}
func mustFinish(t *testing.T, tm *pilosa.TransactionManager, id string) pilosa.Transaction {
t.Helper()
trns, err := tm.Finish(id)
if err != nil {
t.Errorf("finishing transaction: %v", err)
}
return trns
}
func mustGet(t *testing.T, tm *pilosa.TransactionManager, id string) pilosa.Transaction {
t.Helper()
trns, err := tm.Get(id)
if err != nil {
t.Errorf("getting transaction %s: %v", id, err)
}
return trns
}
func mustList(t *testing.T, tm *pilosa.TransactionManager) map[string]pilosa.Transaction {
t.Helper()
trnsMap, err := tm.List()
if err != nil {
t.Errorf("getting transaction list: %v", err)
}
return trnsMap
}
// compareTransactions errors describing how the
// transactions differ (if at all). The deadlines need only be close
// (within 3ms).
func compareTransactions(t *testing.T, trns1, trns2 pilosa.Transaction) {
t.Helper()
if trns1.ID != trns2.ID {
t.Errorf("IDs differ:\n%+v\n%+v", trns1, trns2)
}
if trns1.Active != trns2.Active {
t.Errorf("Actives differ:\n%+v\n%+v", trns1, trns2)
}
if trns1.Exclusive != trns2.Exclusive {
t.Errorf("Exclusives differ:\n%+v\n%+v", trns1, trns2)
}
if trns1.Timeout != trns2.Timeout {
t.Errorf("Timeouts differ:\n%+v\n%+v", trns1, trns2)
}
diff := trns1.Deadline.Sub(trns2.Deadline)
if diff > time.Millisecond*3 || diff < time.Millisecond*-3 {
t.Errorf("Deadlines differ by %v:\n%+v\n%+v", diff, trns1, trns2)
}
if trns1.Stats != trns2.Stats {
t.Errorf("Stats differ:\n%+v\n%+v", trns1, trns2)
}
}
func TestInMemTransactionStore(t *testing.T) {
ims := pilosa.NewInMemTransactionStore()
err := ims.Put(pilosa.Transaction{ID: "blah", Timeout: time.Second})
if err != nil {
t.Fatalf("adding blah: %v", err)
}
trns, err := ims.Get("blah")
if err != nil {
t.Fatalf("getting blah: %v", err)
}
if trns.ID != "blah" || trns.Timeout != time.Second {
t.Fatalf("unexpected transaction for blah: %+v", t)
}
trns, err = ims.Get("nope")
if err != pilosa.ErrTransactionNotFound {
t.Fatalf("unexpected error: %v", err)
}
l, err := ims.List()
if err != nil {
t.Fatalf("listing transactions: %v", err)
}
if len(l) != 1 {
t.Errorf("unexpected number of transactions: %d", len(l))
}
if l["blah"].ID != "blah" || l["blah"].Timeout != time.Second {
t.Errorf("unexpected transaction at blah: %+v", l["blah"])
}
}