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
68 changes: 41 additions & 27 deletions acl/object.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,53 +44,67 @@ type aclObject struct {
store list.Storage

list.AclList
// ready is closed by the first consensus event, which leaves consErr set when the object failed to load
ready chan struct{}
loaded bool
consErr error

lastUsage atomic.Time

mu sync.Mutex
}

// AddConsensusRecords builds the list from the first event and adds the records of the later ones.
// A watch can deliver more events after an error, which only finishes the load once.
func (a *aclObject) AddConsensusRecords(recs []*consensusproto.RawRecordWithId) {
a.mu.Lock()
defer a.mu.Unlock()
slices.Reverse(recs)
if a.store == nil {
defer close(a.ready)
if a.store, a.consErr = list.NewInMemoryStorage(a.id, recs); a.consErr != nil {
return
}
verifier := recordverifier.AcceptorVerifier(recordverifier.NewValidateFull())
if networkId := a.aclService.nodeConf.Configuration().NetworkId; networkId != "" {
netKey, err := crypto.DecodeNetworkId(networkId)
if err != nil {
a.consErr = fmt.Errorf("invalid networkId: %w", err)
return
}
verifier = recordverifier.New(netKey)
}
if a.AclList, a.consErr = list.BuildAclListWithIdentity(a.aclService.accountService.Account(), a.store, verifier); a.consErr != nil {
return
}
} else {
a.Lock()
defer a.Unlock()
if err := a.AddRawRecords(recs); err != nil {
log.Warn("unable to add consensus records", zap.Error(err), zap.String("spaceId", a.id))
return
if !a.loaded {
a.finishLoad(a.build(recs))
return
}
if a.consErr != nil {
// the object failed to load and is being dropped
return
}
a.Lock()
defer a.Unlock()
if err := a.AddRawRecords(recs); err != nil {
log.Warn("unable to add consensus records", zap.Error(err), zap.String("spaceId", a.id))
}
}

func (a *aclObject) build(recs []*consensusproto.RawRecordWithId) (err error) {
if a.store, err = list.NewInMemoryStorage(a.id, recs); err != nil {
return err
}
verifier := recordverifier.AcceptorVerifier(recordverifier.NewValidateFull())
if networkId := a.aclService.nodeConf.Configuration().NetworkId; networkId != "" {
netKey, err := crypto.DecodeNetworkId(networkId)
if err != nil {
return fmt.Errorf("invalid networkId: %w", err)
}
verifier = recordverifier.New(netKey)
}
a.AclList, err = list.BuildAclListWithIdentity(a.aclService.accountService.Account(), a.store, verifier)
return err
}

// finishLoad ends the wait in newAclObject with err; it is called once, under mu
func (a *aclObject) finishLoad(err error) {
a.loaded = true
a.consErr = err
close(a.ready)
}

func (a *aclObject) AddConsensusError(err error) {
a.mu.Lock()
defer a.mu.Unlock()
if a.store == nil {
a.consErr = err
close(a.ready)
if !a.loaded {
a.finishLoad(err)
} else {
log.Warn("got consensus error", zap.Error(err))
log.Warn("got consensus error", zap.Error(err), zap.String("spaceId", a.id))
}
}

Expand Down
72 changes: 72 additions & 0 deletions acl/object_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
package acl

import (
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.uber.org/mock/gomock"

"github.com/anyproto/any-sync/commonspace/object/accountdata"
"github.com/anyproto/any-sync/commonspace/object/acl/list"
"github.com/anyproto/any-sync/consensus/consensusclient"
"github.com/anyproto/any-sync/consensus/consensusproto"
"github.com/anyproto/any-sync/consensus/consensusproto/consensuserr"
)

// The consensus client can deliver more than one event for a watch before the object is dropped:
// it hands a whole batch to the watcher, and a watch restored after a reconnect can report the same
// missing log twice. Only the first event decides how the object loads.
func TestAclObject_EventsAfterTheFirst(t *testing.T) {
ownerKeys, err := accountdata.NewRandom()
require.NoError(t, err)
const spaceId = "spaceId"
ownerAcl, err := list.NewInMemoryDerivedAcl(spaceId, ownerKeys)
require.NoError(t, err)
root := func() []*consensusproto.RawRecordWithId {
return []*consensusproto.RawRecordWithId{ownerAcl.Root()}
}

// load creates the object while the watch delivers events, and waits until all of them are handled
load := func(t *testing.T, events ...func(w consensusclient.Watcher)) (*aclObject, error) {
fx := newFixture(t)
defer fx.finish(t)
handled := make(chan struct{})
fx.consCl.EXPECT().Watch(spaceId, gomock.Any()).DoAndReturn(func(_ string, w consensusclient.Watcher) error {
go func() {
defer close(handled)
for _, event := range events {
event(w)
}
}()
return nil
})
fx.consCl.EXPECT().UnWatch(spaceId).AnyTimes()
obj, err := fx.AclService.(*aclService).newAclObject(ctx, spaceId)
<-handled
return obj, err
}
notFound := func(w consensusclient.Watcher) { w.AddConsensusError(consensuserr.ErrLogNotFound) }
records := func(w consensusclient.Watcher) { w.AddConsensusRecords(root()) }

t.Run("records after an error are dropped", func(t *testing.T) {
_, err := load(t, notFound, records)
assert.ErrorIs(t, err, consensuserr.ErrLogNotFound)
})
t.Run("a second error is dropped", func(t *testing.T) {
_, err := load(t, notFound, notFound)
assert.ErrorIs(t, err, consensuserr.ErrLogNotFound)
})
t.Run("records after records that do not build a list are dropped", func(t *testing.T) {
broken := func(w consensusclient.Watcher) {
w.AddConsensusRecords([]*consensusproto.RawRecordWithId{{Id: "broken", Payload: []byte("broken")}})
}
_, err := load(t, broken, records)
assert.Error(t, err)
})
t.Run("an error after the records leaves the object loaded", func(t *testing.T) {
obj, err := load(t, records, notFound)
require.NoError(t, err)
assert.Equal(t, ownerAcl.Id(), obj.Id())
})
}
Loading