From 38cec6f20e76b3e51845a94a269c70147b484f40 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 15 Apr 2020 17:04:14 -0500 Subject: [PATCH] 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. --- transaction.go | 359 ++++++++++++++++++++++++++++++++++++++++++++ transaction_test.go | 264 ++++++++++++++++++++++++++++++++ 2 files changed, 623 insertions(+) create mode 100644 transaction.go create mode 100644 transaction_test.go diff --git a/transaction.go b/transaction.go new file mode 100644 index 000000000..a38f30848 --- /dev/null +++ b/transaction.go @@ -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") diff --git a/transaction_test.go b/transaction_test.go new file mode 100644 index 000000000..aa64ba446 --- /dev/null +++ b/transaction_test.go @@ -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"]) + } + +}