diff --git a/http/client.go b/http/client.go index b18e5a382..0708c9c42 100644 --- a/http/client.go +++ b/http/client.go @@ -1248,8 +1248,8 @@ func (c *InternalClient) Transactions(ctx context.Context) (map[string]pilosa.Tr return trnsMap, errors.Wrap(err, "executing request") } defer func() { - io.Copy(ioutil.Discard, resp.Body) - resp.Body.Close() + _, _ = io.Copy(ioutil.Discard, resp.Body) + _ = resp.Body.Close() }() tmpTrnsMap := make(map[string]*pilosa.Transaction) err = json.NewDecoder(resp.Body).Decode(&tmpTrnsMap) @@ -1292,8 +1292,8 @@ func (c *InternalClient) StartTransaction(ctx context.Context, id string, timeou return pilosa.Transaction{}, errors.Wrap(err, "executing request") } defer func() { - io.Copy(ioutil.Discard, resp.Body) - resp.Body.Close() + _, _ = io.Copy(ioutil.Discard, resp.Body) + _ = resp.Body.Close() }() err = json.NewDecoder(resp.Body).Decode(&tr) if err != nil { @@ -1325,8 +1325,8 @@ func (c *InternalClient) FinishTransaction(ctx context.Context, id string) (pilo return pilosa.Transaction{}, errors.Wrap(err, "executing request") } defer func() { - io.Copy(ioutil.Discard, resp.Body) - resp.Body.Close() + _, _ = io.Copy(ioutil.Discard, resp.Body) + _ = resp.Body.Close() }() tr := &TransactionResponse{Transaction: &pilosa.Transaction{}} err = json.NewDecoder(resp.Body).Decode(&tr) @@ -1361,8 +1361,8 @@ func (c *InternalClient) GetTransaction(ctx context.Context, id string) (pilosa. return pilosa.Transaction{}, errors.Wrap(err, "executing request") } defer func() { - io.Copy(ioutil.Discard, resp.Body) - resp.Body.Close() + _, _ = io.Copy(ioutil.Discard, resp.Body) + _ = resp.Body.Close() }() tr := &TransactionResponse{Transaction: &pilosa.Transaction{}} err = json.NewDecoder(resp.Body).Decode(&tr) diff --git a/http/client_test.go b/http/client_test.go index 87d820150..7924106c9 100644 --- a/http/client_test.go +++ b/http/client_test.go @@ -1381,7 +1381,6 @@ func TestClientTransactions(t *testing.T) { !strings.Contains(err.Error(), pilosa.ErrNodeNotCoordinator.Error()) { t.Fatalf("unexpected error starting on non-coordinator: %v", err) } else { - expDeadline = time.Now().Add(time.Minute) test.CompareTransactions(t, pilosa.Transaction{}, trns) diff --git a/server.go b/server.go index d4f59996c..143e75e2b 100644 --- a/server.go +++ b/server.go @@ -1050,13 +1050,19 @@ func (srv *Server) StartTransaction(ctx context.Context, id string, timeout time }) if err != nil { // try to clean up, but ignore errors - srv.holder.FinishTransaction(ctx, id) - srv.SendSync( + _, errLocal := srv.holder.FinishTransaction(ctx, id) + errBroadcast := srv.SendSync( &TransactionMessage{ Action: TRANSACTION_FINISH, Transaction: trns, }, ) + if errLocal != nil || errBroadcast != nil { + srv.logger.Printf("error(s) while trying to clean up transaction which failed to start, local: %v, broadcast: %v", + errLocal, + errBroadcast, + ) + } return trns, errors.Wrap(err, "broadcasting transaction start") } return trns, nil diff --git a/test/transaction.go b/test/transaction.go index 3547ad856..a09b9a73d 100644 --- a/test/transaction.go +++ b/test/transaction.go @@ -7,9 +7,11 @@ import ( "github.com/pilosa/pilosa/v2" ) +const deadlineSkew = time.Millisecond * 10 + // CompareTransactions errors describing how the // transactions differ (if at all). The deadlines need only be close -// (within 3ms). +// (within deadlineSkew). func CompareTransactions(t *testing.T, trns1, trns2 pilosa.Transaction) { t.Helper() if trns1.ID != trns2.ID { @@ -27,7 +29,7 @@ func CompareTransactions(t *testing.T, trns1, trns2 pilosa.Transaction) { diff := trns1.Deadline.Sub(trns2.Deadline) - if diff > time.Millisecond*3 || diff < time.Millisecond*-3 { + if diff > deadlineSkew || diff < -deadlineSkew { t.Errorf("Deadlines differ by %v:\n%+v\n%+v", diff, trns1, trns2) } if trns1.Stats != trns2.Stats { diff --git a/transaction.go b/transaction.go index 2d890d905..826626883 100644 --- a/transaction.go +++ b/transaction.go @@ -376,6 +376,8 @@ func CompareTransactions(t1, t2 Transaction) error { return nil } +const RFC3339NanoNoZone = "2006-01-02T15:04:05.999999999" + func (trns *Transaction) UnmarshalJSON(b []byte) error { tmp := &struct { ID string `json:"id"` @@ -410,7 +412,7 @@ func (trns *Transaction) UnmarshalJSON(b []byte) error { } if tmp.Deadline != "" { - trns.Deadline, err = time.Parse(time.RFC3339Nano, tmp.Deadline) + trns.Deadline, err = time.ParseInLocation(RFC3339NanoNoZone, tmp.Deadline, time.UTC) } return errors.Wrap(err, "parsing deadline") } @@ -427,6 +429,6 @@ func (trns *Transaction) MarshalJSON() ([]byte, error) { Active: trns.Active, Exclusive: trns.Exclusive, Timeout: trns.Timeout.String(), - Deadline: trns.Deadline.Format(time.RFC3339Nano), + Deadline: trns.Deadline.In(time.UTC).Format(RFC3339NanoNoZone), }) } diff --git a/transaction.md b/transaction.md index 389000a23..dd72d4b34 100644 --- a/transaction.md +++ b/transaction.md @@ -162,7 +162,7 @@ goes through API (and is passed directly to Server). (unimplemented) - [x] implement api layer and cluster logic, startup, etc. - [ ] add new cluster state to explicitly reject certain requests during exclusive transaction? - [x] implement HTTP layer -- [x] implement transaction id in header +- [ ] implement transaction id in header - [x] propagate context - [ ] implement and use persistent transaction store rather than inmem. - [ ] update go-pilosa/gpexp to actually USE transactions