-
Notifications
You must be signed in to change notification settings - Fork 70
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
refactor(kafkatopic): add KafkaTopic repository [skip ci]
- Loading branch information
Showing
16 changed files
with
769 additions
and
708 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,44 @@ | ||
package kafkatopicrepository | ||
|
||
import ( | ||
"context" | ||
"time" | ||
|
||
"github.com/aiven/aiven-go-client" | ||
"github.com/avast/retry-go" | ||
) | ||
|
||
const createRetryDelay = 10 * time.Second | ||
|
||
// Create creates topic. | ||
// First checks if topic does not exist for the safety | ||
// Then calls creates topic. | ||
func (rep *repository) Create(ctx context.Context, project, service string, req aiven.CreateKafkaTopicRequest) error { | ||
// aiven.KafkaTopics.Create() function may return 501 on create | ||
// Second call might say that topic already exists, and we have retries in aiven client | ||
// So to be sure, better check it before create | ||
err := rep.exists(ctx, project, service, req.TopicName, true) | ||
if err == nil { | ||
return errAlreadyExists | ||
} | ||
|
||
// If this is not errNotFound, then something happened | ||
if err != errNotFound { | ||
return err | ||
} | ||
|
||
return retry.Do(func() error { | ||
err := rep.client.Create(project, service, req) | ||
if err == nil || aiven.IsAlreadyExists(err) { | ||
return nil | ||
} | ||
|
||
// If some brokers are offline while the request is being executed | ||
// the operation may fail. | ||
_, ok := err.(aiven.Error) | ||
if ok { | ||
return err | ||
} | ||
return retry.Unrecoverable(err) | ||
}, retry.Context(ctx), retry.Delay(createRetryDelay)) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,68 @@ | ||
package kafkatopicrepository | ||
|
||
import ( | ||
"context" | ||
"sync" | ||
"sync/atomic" | ||
"testing" | ||
|
||
"github.com/aiven/aiven-go-client" | ||
"github.com/stretchr/testify/assert" | ||
) | ||
|
||
// TestCreateConflict tests that one goroutine out of 100 creates topic, while others get errAlreadyExists | ||
func TestCreateConflict(t *testing.T) { | ||
client := &fakeTopicClient{} | ||
rep := newRepository(client) | ||
ctx := context.Background() | ||
|
||
var conflictErr int32 | ||
var wg sync.WaitGroup | ||
for i := 0; i < 100; i++ { | ||
wg.Add(1) | ||
go func() { | ||
err := rep.Create(ctx, "foo", "bar", aiven.CreateKafkaTopicRequest{TopicName: "baz"}) | ||
if err == errAlreadyExists { | ||
atomic.AddInt32(&conflictErr, 1) | ||
} | ||
wg.Done() | ||
}() | ||
} | ||
wg.Wait() | ||
assert.EqualValues(t, 99, conflictErr) | ||
assert.EqualValues(t, 1, client.createCalled) | ||
assert.EqualValues(t, 1, client.v1ListCalled) | ||
assert.EqualValues(t, 0, client.v2ListCalled) | ||
assert.True(t, rep.seenServices["foo/bar"]) | ||
assert.True(t, rep.seenTopics["foo/bar/baz"]) | ||
} | ||
|
||
func TestCreateRecreateMissing(t *testing.T) { | ||
client := &fakeTopicClient{} | ||
rep := newRepository(client) | ||
ctx := context.Background() | ||
|
||
// Creates topic | ||
err := rep.Create(ctx, "foo", "bar", aiven.CreateKafkaTopicRequest{TopicName: "baz"}) | ||
assert.NoError(t, err) | ||
assert.EqualValues(t, 1, client.createCalled) | ||
assert.EqualValues(t, 1, client.v1ListCalled) | ||
assert.EqualValues(t, 0, client.v2ListCalled) | ||
assert.True(t, rep.seenServices["foo/bar"]) | ||
assert.True(t, rep.seenTopics["foo/bar/baz"]) | ||
|
||
// Forgets the topic, like if it's missing | ||
err = rep.forgetTopic("foo", "bar", "baz") | ||
assert.NoError(t, err) | ||
assert.True(t, rep.seenServices["foo/bar"]) | ||
assert.False(t, rep.seenTopics["foo/bar/baz"]) // not cached, missing | ||
|
||
// Recreates topic | ||
err = rep.Create(ctx, "foo", "bar", aiven.CreateKafkaTopicRequest{TopicName: "baz"}) | ||
assert.NoError(t, err) | ||
assert.EqualValues(t, 2, client.createCalled) // Updated | ||
assert.EqualValues(t, 1, client.v1ListCalled) | ||
assert.EqualValues(t, 0, client.v2ListCalled) | ||
assert.True(t, rep.seenServices["foo/bar"]) | ||
assert.True(t, rep.seenTopics["foo/bar/baz"]) // cached again | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,18 @@ | ||
package kafkatopicrepository | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/aiven/aiven-go-client" | ||
) | ||
|
||
func (rep *repository) Delete(_ context.Context, project, service, topic string) error { | ||
rep.Lock() | ||
rep.seenTopics[newKey(project, service, topic)] = false | ||
rep.Unlock() | ||
err := rep.client.Delete(project, service, topic) | ||
if !(err == nil || aiven.IsNotFound(err)) { | ||
return err | ||
} | ||
return nil | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,163 @@ | ||
package kafkatopicrepository | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
|
||
"github.com/aiven/aiven-go-client" | ||
"github.com/avast/retry-go" | ||
"github.com/samber/lo" | ||
) | ||
|
||
func (rep *repository) Read(ctx context.Context, project, service, topic string) (*aiven.KafkaTopic, error) { | ||
// We have quick methods to determine that topic does not exist | ||
err := rep.exists(ctx, project, service, topic, false) | ||
if err != nil { | ||
return nil, err | ||
} | ||
|
||
// Adds request to the queue | ||
c := make(chan *response, 1) | ||
r := &request{ | ||
project: project, | ||
service: service, | ||
topic: topic, | ||
rsp: c, | ||
} | ||
rep.Lock() | ||
rep.queue = append(rep.queue, r) | ||
rep.Unlock() | ||
|
||
// Waits response from the channel | ||
// Or exits on context done | ||
select { | ||
case <-ctx.Done(): | ||
return nil, ctx.Err() | ||
case rsp := <-c: | ||
close(c) | ||
return rsp.topic, rsp.err | ||
} | ||
} | ||
|
||
// exists returns nil if topic exists, or errNotFound if doesn't: | ||
// 1. checks repository.seenTopics for known topics | ||
// 2. calls v1List for the remote state for the given service and marks it in repository.seenServices | ||
// 3. saves topic names to repository.seenTopics, so its result can be reused | ||
// 4. when acquire true, then saves topic to repository.seenTopics (for creating) | ||
// todo: use context with the new client | ||
func (rep *repository) exists(_ context.Context, project, service, topic string, acquire bool) error { | ||
rep.Lock() | ||
defer rep.Unlock() | ||
// Checks repository.seenTopics. | ||
// If it has been just created, it is not available in v1List. | ||
// So calling it first doesn't make any sense | ||
serviceKey := newKey(project, service) | ||
topicKey := newKey(serviceKey, topic) | ||
if rep.seenTopics[topicKey] { | ||
return nil | ||
} | ||
|
||
// Goes for v1List | ||
if !rep.seenServices[serviceKey] { | ||
list, err := rep.client.List(project, service) | ||
if err != nil { | ||
return err | ||
} | ||
|
||
// Marks seen all the topics | ||
for _, t := range list { | ||
rep.seenTopics[newKey(serviceKey, t.TopicName)] = true | ||
} | ||
|
||
// Service is seen too. It never goes here again | ||
rep.seenServices[serviceKey] = true | ||
} | ||
|
||
// Checks updated list | ||
if rep.seenTopics[topicKey] { | ||
return nil | ||
} | ||
|
||
// Create functions run in parallel need to lock the name before create | ||
// Otherwise they may run into conflict | ||
if acquire { | ||
rep.seenTopics[topicKey] = true | ||
} | ||
|
||
// v1List doesn't contain the topic | ||
return errNotFound | ||
} | ||
|
||
// fetch fetches requested topics configuration | ||
// 1. groups topics by service | ||
// 2. requests topics (in chunks) | ||
// Warning: if we call V2List with at least one "not found" topic, it will return 404 for all topics | ||
// Should be certain that all topics in queue do exist. Call repository.exists first to do so | ||
func (rep *repository) fetch(queue map[string]*request) { | ||
// Groups topics by service | ||
byService := make(map[string][]*request, 0) | ||
for i := range queue { | ||
r := queue[i] | ||
key := newKey(r.project, r.service) | ||
byService[key] = append(byService[key], r) | ||
} | ||
|
||
// Fetches topics configuration | ||
for _, reqs := range byService { | ||
topicNames := make([]string, 0, len(reqs)) | ||
for _, r := range reqs { | ||
topicNames = append(topicNames, r.topic) | ||
} | ||
|
||
// Topics are grouped by service | ||
// We can share this values | ||
project := reqs[0].project | ||
service := reqs[0].service | ||
|
||
// Slices topic names by repository.v2ListBatchSize | ||
// because V2List has a limit | ||
for _, chunk := range lo.Chunk(topicNames, rep.v2ListBatchSize) { | ||
// V2List() and Get() do not get info immediately | ||
// Some retries should be applied if result is not equal to requested values | ||
var list []*aiven.KafkaTopic | ||
err := retry.Do(func() error { | ||
rspList, err := rep.client.V2List(project, service, chunk) | ||
|
||
// 404 means that there is "not found" on the list | ||
// But repository.exists should have checked these, so now this is a fail | ||
if aiven.IsNotFound(err) { | ||
return retry.Unrecoverable(fmt.Errorf("topic list has changed")) | ||
} | ||
|
||
// Something else happened | ||
// We have retries in the client, so this is bad | ||
if err != nil { | ||
return retry.Unrecoverable(err) | ||
} | ||
|
||
// This is an old cache, we need to retry it until succeed | ||
if len(rspList) != len(chunk) { | ||
return fmt.Errorf("got %d topics, expected %d. Retrying", len(rspList), len(chunk)) | ||
} | ||
|
||
list = rspList | ||
return nil | ||
}, retry.Delay(rep.v2ListRetryDelay)) | ||
|
||
if err != nil { | ||
// Send errors | ||
// Flattens error to a string, because it might go really completed for testing | ||
err = fmt.Errorf("topic read error: %s", err) | ||
for _, r := range reqs { | ||
r.send(nil, err) | ||
} | ||
continue | ||
} | ||
|
||
// Sends topics | ||
for _, t := range list { | ||
queue[newKey(project, service, t.TopicName)].send(t, nil) | ||
} | ||
} | ||
} | ||
} |
Oops, something went wrong.