Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 29 additions & 4 deletions pkg/service/leader/leader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,24 +2,25 @@ package leader_test

import (
"context"
"errors"
"testing"
"testing/synctest"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

"go.atoms.co/lib/chanx"
"go.atoms.co/lib/testing/assertx"
"go.atoms.co/splitter/lib/service/location"
"go.atoms.co/splitter/lib/service/session"
"go.atoms.co/lib/testing/assertx"
splitterpb "go.atoms.co/splitter/pb"
splitterprivatepb "go.atoms.co/splitter/pb/private"
"go.atoms.co/splitter/pkg/core"
"go.atoms.co/splitter/pkg/model"
"go.atoms.co/splitter/pkg/service/leader"
"go.atoms.co/splitter/pkg/storage"
"go.atoms.co/splitter/pkg/storage/memory"
splitterprivatepb "go.atoms.co/splitter/pb/private"
splitterpb "go.atoms.co/splitter/pb"
"go.atoms.co/lib/chanx"
)

const (
Expand Down Expand Up @@ -190,6 +191,22 @@ func TestLeader_Operations(t *testing.T) {
assert.Len(t, snap.GetSnapshot().GetTenants(), 2)
}

func TestLeader_DoesNotAcknowledgeFailedUpdate(t *testing.T) {
ctx := context.Background()
loc := location.New("centralus", "splitter-0")
db := failingUpdateStorage{Storage: memory.New()}

l := leader.New(ctx, loc, db, leader.WithFastActivation())
defer l.Close()
<-l.Initialized().Closed()

response, err := l.Handle(ctx, leader.NewHandleTenantRequest(&splitterprivatepb.TenantRequest{
Req: &splitterprivatepb.TenantRequest_New{New: &splitterpb.NewTenantRequest{Name: string(tenant1)}},
}))
require.ErrorIs(t, err, model.ErrNotOwned)
require.Nil(t, response)
}

func TestLeader_HandleUpdate(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
ctx := context.Background()
Expand Down Expand Up @@ -283,6 +300,14 @@ func setup(t *testing.T, ctx context.Context, services ...model.Service) storage
return db
}

type failingUpdateStorage struct {
storage.Storage
}

func (failingUpdateStorage) Update(context.Context, core.Update) error {
return errors.New("apply failed")
}

func setupWithDomains(t *testing.T, ctx context.Context, service model.Service, domains ...model.Domain) storage.Storage {
db := memory.New()

Expand Down
17 changes: 8 additions & 9 deletions pkg/service/leader/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,14 +5,14 @@ import (
"fmt"
"time"

"go.atoms.co/lib/log"
"go.atoms.co/iox"
"go.atoms.co/lib/log"
"go.atoms.co/lib/workqueue"
splitterpb "go.atoms.co/splitter/pb"
splitterprivatepb "go.atoms.co/splitter/pb/private"
"go.atoms.co/splitter/pkg/core"
"go.atoms.co/splitter/pkg/model"
"go.atoms.co/splitter/pkg/storage"
splitterprivatepb "go.atoms.co/splitter/pb/private"
splitterpb "go.atoms.co/splitter/pb"
)

const (
Expand Down Expand Up @@ -612,8 +612,6 @@ func (w *Writer) updateAsync(ctx context.Context, upd core.Update) iox.AsyncClos
func (w *Writer) applyUpdateAsync(ctx context.Context, upd core.Update) iox.AsyncCloser {
done := iox.NewAsyncCloser()
w.pool.Chan() <- func() {
defer done.Close()

// Perform I/O async. If it fails, escalate.

if err := w.db.Update(ctx, upd); err != nil {
Expand All @@ -627,6 +625,7 @@ func (w *Writer) applyUpdateAsync(ctx context.Context, upd core.Update) iox.Asyn
case <-w.Closed():
}

done.Close()
}
return done
}
Expand Down Expand Up @@ -654,8 +653,6 @@ func (w *Writer) deleteAsync(ctx context.Context, del core.Delete) iox.AsyncClos

done := iox.NewAsyncCloser()
w.pool.Chan() <- func() {
defer done.Close()

// Perform I/O async. If it fails, escalate.

if err := w.db.Delete(ctx, del); err != nil {
Expand All @@ -668,6 +665,8 @@ func (w *Writer) deleteAsync(ctx context.Context, del core.Delete) iox.AsyncClos
case w.del <- del:
case <-w.Closed():
}

done.Close()
}
return done
}
Expand All @@ -677,8 +676,6 @@ func (w *Writer) restoreAsync(ctx context.Context, res core.Restore) iox.AsyncCl

done := iox.NewAsyncCloser()
w.pool.Chan() <- func() {
defer done.Close()

// Perform I/O async. If it fails, escalate.

if err := w.db.Restore(ctx, res); err != nil {
Expand All @@ -691,6 +688,8 @@ func (w *Writer) restoreAsync(ctx context.Context, res core.Restore) iox.AsyncCl
case w.res <- res:
case <-w.Closed():
}

done.Close()
}
return done
}
Expand Down