mirror of
https://github.com/zitadel/zitadel.git
synced 2025-05-30 19:28:41 +00:00
try with goroutines
This commit is contained in:
parent
55e5e82dbc
commit
35ce026651
@ -6,6 +6,7 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"regexp"
|
"regexp"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
"sync"
|
||||||
|
|
||||||
"github.com/caos/logging"
|
"github.com/caos/logging"
|
||||||
caos_errs "github.com/caos/zitadel/internal/errors"
|
caos_errs "github.com/caos/zitadel/internal/errors"
|
||||||
@ -109,45 +110,59 @@ func (db *CRDB) Push(ctx context.Context, events ...*repository.Event) error {
|
|||||||
err := crdb.ExecuteTx(ctx, db.client, nil, func(tx *sql.Tx) error {
|
err := crdb.ExecuteTx(ctx, db.client, nil, func(tx *sql.Tx) error {
|
||||||
stmt, err := tx.PrepareContext(ctx, crdbInsert)
|
stmt, err := tx.PrepareContext(ctx, crdbInsert)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
tx.Rollback()
|
|
||||||
logging.Log("SQL-3to5p").WithError(err).Warn("prepare failed")
|
logging.Log("SQL-3to5p").WithError(err).Warn("prepare failed")
|
||||||
return caos_errs.ThrowInternal(err, "SQL-OdXRE", "prepare failed")
|
return caos_errs.ThrowInternal(err, "SQL-OdXRE", "prepare failed")
|
||||||
}
|
}
|
||||||
|
wg := sync.WaitGroup{}
|
||||||
|
errs := make(chan error, len(events))
|
||||||
|
|
||||||
for _, event := range events {
|
for _, event := range events {
|
||||||
previousSequence := Sequence(event.PreviousSequence)
|
wg.Add(1)
|
||||||
if event.PreviousEvent != nil {
|
go func(event *repository.Event) {
|
||||||
previousSequence = Sequence(event.PreviousEvent.Sequence)
|
defer wg.Done()
|
||||||
}
|
previousSequence := Sequence(event.PreviousSequence)
|
||||||
err = stmt.QueryRowContext(ctx,
|
if event.PreviousEvent != nil {
|
||||||
event.Type,
|
if event.PreviousEvent.AggregateType != event.AggregateType || event.PreviousEvent.AggregateID != event.AggregateID {
|
||||||
event.AggregateType,
|
errs <- caos_errs.ThrowPreconditionFailed(nil, "SQL-J55uR", "aggregate of linked events unequal")
|
||||||
event.AggregateID,
|
return
|
||||||
event.Version,
|
}
|
||||||
&sql.NullTime{
|
previousSequence = Sequence(event.PreviousEvent.Sequence)
|
||||||
Time: event.CreationDate,
|
}
|
||||||
Valid: !event.CreationDate.IsZero(),
|
err = stmt.QueryRowContext(ctx,
|
||||||
},
|
event.Type,
|
||||||
Data(event.Data),
|
event.AggregateType,
|
||||||
event.EditorUser,
|
event.AggregateID,
|
||||||
event.EditorService,
|
event.Version,
|
||||||
event.ResourceOwner,
|
&sql.NullTime{
|
||||||
previousSequence,
|
Time: event.CreationDate,
|
||||||
event.CheckPreviousSequence,
|
Valid: !event.CreationDate.IsZero(),
|
||||||
).Scan(&event.ID, &event.Sequence, &previousSequence, &event.CreationDate)
|
},
|
||||||
|
Data(event.Data),
|
||||||
|
event.EditorUser,
|
||||||
|
event.EditorService,
|
||||||
|
event.ResourceOwner,
|
||||||
|
previousSequence,
|
||||||
|
event.CheckPreviousSequence,
|
||||||
|
).Scan(&event.ID, &event.Sequence, &previousSequence, &event.CreationDate)
|
||||||
|
|
||||||
event.PreviousSequence = uint64(previousSequence)
|
event.PreviousSequence = uint64(previousSequence)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
tx.Rollback()
|
logging.LogWithFields("SQL-IP3js",
|
||||||
|
"aggregate", event.AggregateType,
|
||||||
logging.LogWithFields("SQL-IP3js",
|
"aggregateId", event.AggregateID,
|
||||||
"aggregate", event.AggregateType,
|
"aggregateType", event.AggregateType,
|
||||||
"aggregateId", event.AggregateID,
|
"eventType", event.Type).WithError(err).Info("query failed")
|
||||||
"aggregateType", event.AggregateType,
|
errs <- caos_errs.ThrowInternal(err, "SQL-SBP37", "unable to create event")
|
||||||
"eventType", event.Type).WithError(err).Info("query failed")
|
}
|
||||||
return caos_errs.ThrowInternal(err, "SQL-SBP37", "unable to create event")
|
}(event)
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
wg.Wait()
|
||||||
|
close(errs)
|
||||||
|
for err := range errs {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
})
|
})
|
||||||
if err != nil && !errors.Is(err, &caos_errs.CaosError{}) {
|
if err != nil && !errors.Is(err, &caos_errs.CaosError{}) {
|
||||||
|
@ -2,6 +2,7 @@ package sql
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/caos/zitadel/internal/eventstore/v2/repository"
|
"github.com/caos/zitadel/internal/eventstore/v2/repository"
|
||||||
@ -265,6 +266,7 @@ func TestCRDB_columnName(t *testing.T) {
|
|||||||
|
|
||||||
func TestCRDB_Push_OneAggregate(t *testing.T) {
|
func TestCRDB_Push_OneAggregate(t *testing.T) {
|
||||||
type args struct {
|
type args struct {
|
||||||
|
ctx context.Context
|
||||||
events []*repository.Event
|
events []*repository.Event
|
||||||
}
|
}
|
||||||
type eventsRes struct {
|
type eventsRes struct {
|
||||||
@ -281,23 +283,10 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
args args
|
args args
|
||||||
res res
|
res res
|
||||||
}{
|
}{
|
||||||
{
|
|
||||||
name: "push no events",
|
|
||||||
args: args{
|
|
||||||
events: []*repository.Event{},
|
|
||||||
},
|
|
||||||
res: res{
|
|
||||||
wantErr: false,
|
|
||||||
eventsRes: eventsRes{
|
|
||||||
pushedEventsCount: 0,
|
|
||||||
aggID: []string{"0"},
|
|
||||||
aggType: repository.AggregateType(t.Name()),
|
|
||||||
},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
{
|
{
|
||||||
name: "push 1 event with check previous",
|
name: "push 1 event with check previous",
|
||||||
args: args{
|
args: args{
|
||||||
|
ctx: context.Background(),
|
||||||
events: []*repository.Event{
|
events: []*repository.Event{
|
||||||
generateEvent(t, "1", true, 0),
|
generateEvent(t, "1", true, 0),
|
||||||
},
|
},
|
||||||
@ -313,6 +302,7 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
{
|
{
|
||||||
name: "fail push 1 event with check previous wrong sequence",
|
name: "fail push 1 event with check previous wrong sequence",
|
||||||
args: args{
|
args: args{
|
||||||
|
ctx: context.Background(),
|
||||||
events: []*repository.Event{
|
events: []*repository.Event{
|
||||||
generateEvent(t, "2", true, 5),
|
generateEvent(t, "2", true, 5),
|
||||||
},
|
},
|
||||||
@ -329,6 +319,7 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
{
|
{
|
||||||
name: "push 1 event without check previous",
|
name: "push 1 event without check previous",
|
||||||
args: args{
|
args: args{
|
||||||
|
ctx: context.Background(),
|
||||||
events: []*repository.Event{
|
events: []*repository.Event{
|
||||||
generateEvent(t, "3", false, 0),
|
generateEvent(t, "3", false, 0),
|
||||||
},
|
},
|
||||||
@ -345,6 +336,7 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
{
|
{
|
||||||
name: "push 1 event without check previous wrong sequence",
|
name: "push 1 event without check previous wrong sequence",
|
||||||
args: args{
|
args: args{
|
||||||
|
ctx: context.Background(),
|
||||||
events: []*repository.Event{
|
events: []*repository.Event{
|
||||||
generateEvent(t, "4", false, 5),
|
generateEvent(t, "4", false, 5),
|
||||||
},
|
},
|
||||||
@ -361,6 +353,7 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
{
|
{
|
||||||
name: "fail on push two events on agg without linking",
|
name: "fail on push two events on agg without linking",
|
||||||
args: args{
|
args: args{
|
||||||
|
ctx: context.Background(),
|
||||||
events: []*repository.Event{
|
events: []*repository.Event{
|
||||||
generateEvent(t, "5", true, 0),
|
generateEvent(t, "5", true, 0),
|
||||||
generateEvent(t, "5", true, 0),
|
generateEvent(t, "5", true, 0),
|
||||||
@ -378,6 +371,7 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
{
|
{
|
||||||
name: "push two events on agg with linking",
|
name: "push two events on agg with linking",
|
||||||
args: args{
|
args: args{
|
||||||
|
ctx: context.Background(),
|
||||||
events: linkEvents(
|
events: linkEvents(
|
||||||
generateEvent(t, "6", true, 0),
|
generateEvent(t, "6", true, 0),
|
||||||
generateEvent(t, "6", true, 0),
|
generateEvent(t, "6", true, 0),
|
||||||
@ -395,6 +389,7 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
{
|
{
|
||||||
name: "push two events on agg with linking without check previous",
|
name: "push two events on agg with linking without check previous",
|
||||||
args: args{
|
args: args{
|
||||||
|
ctx: context.Background(),
|
||||||
events: linkEvents(
|
events: linkEvents(
|
||||||
generateEvent(t, "7", false, 0),
|
generateEvent(t, "7", false, 0),
|
||||||
generateEvent(t, "7", false, 0),
|
generateEvent(t, "7", false, 0),
|
||||||
@ -412,6 +407,7 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
{
|
{
|
||||||
name: "push two events on agg with linking mixed check previous",
|
name: "push two events on agg with linking mixed check previous",
|
||||||
args: args{
|
args: args{
|
||||||
|
ctx: context.Background(),
|
||||||
events: linkEvents(
|
events: linkEvents(
|
||||||
generateEvent(t, "8", false, 0),
|
generateEvent(t, "8", false, 0),
|
||||||
generateEvent(t, "8", true, 0),
|
generateEvent(t, "8", true, 0),
|
||||||
@ -429,13 +425,30 @@ func TestCRDB_Push_OneAggregate(t *testing.T) {
|
|||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
|
{
|
||||||
|
name: "failed push because context canceled",
|
||||||
|
args: args{
|
||||||
|
ctx: canceledCtx(),
|
||||||
|
events: []*repository.Event{
|
||||||
|
generateEvent(t, "9", true, 0),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
res: res{
|
||||||
|
wantErr: true,
|
||||||
|
eventsRes: eventsRes{
|
||||||
|
pushedEventsCount: 0,
|
||||||
|
aggID: []string{"9"},
|
||||||
|
aggType: repository.AggregateType(t.Name()),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
}
|
}
|
||||||
for _, tt := range tests {
|
for _, tt := range tests {
|
||||||
t.Run(tt.name, func(t *testing.T) {
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
db := &CRDB{
|
db := &CRDB{
|
||||||
client: testCRDBClient,
|
client: testCRDBClient,
|
||||||
}
|
}
|
||||||
if err := db.Push(context.Background(), tt.args.events...); (err != nil) != tt.res.wantErr {
|
if err := db.Push(tt.args.ctx, tt.args.events...); (err != nil) != tt.res.wantErr {
|
||||||
t.Errorf("CRDB.Push() error = %v, wantErr %v", err, tt.res.wantErr)
|
t.Errorf("CRDB.Push() error = %v, wantErr %v", err, tt.res.wantErr)
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -516,14 +529,60 @@ func TestCRDB_Push_MultipleAggregate(t *testing.T) {
|
|||||||
args: args{
|
args: args{
|
||||||
events: linkEvents(
|
events: linkEvents(
|
||||||
generateEvent(t, "104", false, 0),
|
generateEvent(t, "104", false, 0),
|
||||||
generateEvent(t, "104", false, 0),
|
generateEvent(t, "105", false, 0),
|
||||||
),
|
),
|
||||||
},
|
},
|
||||||
res: res{
|
res: res{
|
||||||
wantErr: false,
|
wantErr: true,
|
||||||
eventsRes: eventsRes{
|
eventsRes: eventsRes{
|
||||||
pushedEventsCount: 0,
|
pushedEventsCount: 0,
|
||||||
aggID: []string{"104"},
|
aggID: []string{"104", "105"},
|
||||||
|
aggType: []repository.AggregateType{repository.AggregateType(t.Name())},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "push two aggregates mixed check previous multiple events",
|
||||||
|
args: args{
|
||||||
|
events: combineEventLists(
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "106", true, 0),
|
||||||
|
generateEvent(t, "106", false, 0),
|
||||||
|
generateEvent(t, "106", false, 0),
|
||||||
|
generateEvent(t, "106", true, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "107", false, 0),
|
||||||
|
generateEvent(t, "107", true, 0),
|
||||||
|
generateEvent(t, "107", false, 0),
|
||||||
|
generateEvent(t, "107", true, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "108", true, 0),
|
||||||
|
generateEvent(t, "108", false, 0),
|
||||||
|
generateEvent(t, "108", false, 0),
|
||||||
|
generateEvent(t, "108", true, 0),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "failed push same aggregate in two transactions",
|
||||||
|
args: args{
|
||||||
|
events: combineEventLists(
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "109", true, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "109", true, 0),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
},
|
||||||
|
res: res{
|
||||||
|
wantErr: true,
|
||||||
|
eventsRes: eventsRes{
|
||||||
|
pushedEventsCount: 0,
|
||||||
|
aggID: []string{"109"},
|
||||||
aggType: []repository.AggregateType{repository.AggregateType(t.Name())},
|
aggType: []repository.AggregateType{repository.AggregateType(t.Name())},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
@ -552,8 +611,383 @@ func TestCRDB_Push_MultipleAggregate(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func combineEventLists(firstList []*repository.Event, secondList []*repository.Event) []*repository.Event {
|
func TestCRDB_Push_Parallel(t *testing.T) {
|
||||||
return append(firstList, secondList...)
|
type args struct {
|
||||||
|
events [][]*repository.Event
|
||||||
|
}
|
||||||
|
type eventsRes struct {
|
||||||
|
pushedEventsCount int
|
||||||
|
aggTypes []repository.AggregateType
|
||||||
|
aggIDs []string
|
||||||
|
}
|
||||||
|
type res struct {
|
||||||
|
errCount int
|
||||||
|
eventsRes eventsRes
|
||||||
|
}
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
args args
|
||||||
|
res res
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
name: "clients push different aggregates",
|
||||||
|
args: args{
|
||||||
|
events: [][]*repository.Event{
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "200", false, 0),
|
||||||
|
generateEvent(t, "200", true, 0),
|
||||||
|
generateEvent(t, "200", false, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "201", false, 0),
|
||||||
|
generateEvent(t, "201", true, 0),
|
||||||
|
generateEvent(t, "201", false, 0),
|
||||||
|
),
|
||||||
|
combineEventLists(
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "202", false, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "203", true, 0),
|
||||||
|
generateEvent(t, "203", false, 0),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
res: res{
|
||||||
|
errCount: 0,
|
||||||
|
eventsRes: eventsRes{
|
||||||
|
aggIDs: []string{"200", "201", "202", "203"},
|
||||||
|
pushedEventsCount: 9,
|
||||||
|
aggTypes: []repository.AggregateType{repository.AggregateType(t.Name())},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "clients push same aggregates no check previous",
|
||||||
|
args: args{
|
||||||
|
events: [][]*repository.Event{
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "204", false, 0),
|
||||||
|
generateEvent(t, "204", false, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "204", false, 0),
|
||||||
|
generateEvent(t, "204", false, 0),
|
||||||
|
),
|
||||||
|
combineEventLists(
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "205", false, 0),
|
||||||
|
generateEvent(t, "205", false, 0),
|
||||||
|
generateEvent(t, "205", false, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "206", false, 0),
|
||||||
|
generateEvent(t, "206", false, 0),
|
||||||
|
generateEvent(t, "206", false, 0),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
combineEventLists(
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "204", false, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "205", false, 0),
|
||||||
|
generateEvent(t, "205", false, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "206", false, 0),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
res: res{
|
||||||
|
errCount: 0,
|
||||||
|
eventsRes: eventsRes{
|
||||||
|
aggIDs: []string{"204", "205", "206"},
|
||||||
|
pushedEventsCount: 14,
|
||||||
|
aggTypes: []repository.AggregateType{repository.AggregateType(t.Name())},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "clients push different aggregates one with check previous",
|
||||||
|
args: args{
|
||||||
|
events: [][]*repository.Event{
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "207", false, 0),
|
||||||
|
generateEvent(t, "207", false, 0),
|
||||||
|
generateEvent(t, "207", false, 0),
|
||||||
|
generateEvent(t, "207", false, 0),
|
||||||
|
generateEvent(t, "207", false, 0),
|
||||||
|
generateEvent(t, "207", false, 0),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEvent(t, "208", true, 0),
|
||||||
|
generateEvent(t, "208", true, 0),
|
||||||
|
generateEvent(t, "208", true, 0),
|
||||||
|
generateEvent(t, "208", true, 0),
|
||||||
|
generateEvent(t, "208", true, 0),
|
||||||
|
),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
res: res{
|
||||||
|
errCount: 0,
|
||||||
|
eventsRes: eventsRes{
|
||||||
|
aggIDs: []string{"207", "208"},
|
||||||
|
pushedEventsCount: 11,
|
||||||
|
aggTypes: []repository.AggregateType{repository.AggregateType(t.Name())},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "clients push different aggregates all with check previous on first event fail",
|
||||||
|
args: args{
|
||||||
|
events: [][]*repository.Event{
|
||||||
|
linkEvents(
|
||||||
|
generateEventWithData(t, "210", true, 0, []byte(`{ "transaction": 1 }`)),
|
||||||
|
generateEventWithData(t, "210", false, 0, []byte(`{ "transaction": 1.1 }`)),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEventWithData(t, "210", true, 0, []byte(`{ "transaction": 2 }`)),
|
||||||
|
generateEventWithData(t, "210", false, 0, []byte(`{ "transaction": 2.1 }`)),
|
||||||
|
),
|
||||||
|
linkEvents(
|
||||||
|
generateEventWithData(t, "210", true, 0, []byte(`{ "transaction": 3 }`)),
|
||||||
|
generateEventWithData(t, "210", false, 0, []byte(`{ "transaction": 30.1 }`)),
|
||||||
|
),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
res: res{
|
||||||
|
errCount: 2,
|
||||||
|
eventsRes: eventsRes{
|
||||||
|
aggIDs: []string{"210"},
|
||||||
|
pushedEventsCount: 2,
|
||||||
|
aggTypes: []repository.AggregateType{repository.AggregateType(t.Name())},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
db := &CRDB{
|
||||||
|
client: testCRDBClient,
|
||||||
|
}
|
||||||
|
wg := sync.WaitGroup{}
|
||||||
|
|
||||||
|
errs := make([]error, 0, tt.res.errCount)
|
||||||
|
errsMu := sync.Mutex{}
|
||||||
|
for _, events := range tt.args.events {
|
||||||
|
wg.Add(1)
|
||||||
|
go func(events []*repository.Event) {
|
||||||
|
err := db.Push(context.Background(), events...)
|
||||||
|
if err != nil {
|
||||||
|
errsMu.Lock()
|
||||||
|
errs = append(errs, err)
|
||||||
|
errsMu.Unlock()
|
||||||
|
}
|
||||||
|
|
||||||
|
wg.Done()
|
||||||
|
}(events)
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
|
||||||
|
if len(errs) != tt.res.errCount {
|
||||||
|
t.Errorf("CRDB.Push() error count = %d, wanted err count %d, errs: %v", len(errs), tt.res.errCount, errs)
|
||||||
|
}
|
||||||
|
|
||||||
|
rows, err := testCRDBClient.Query("SELECT event_data FROM eventstore.events where aggregate_type = ANY($1) AND aggregate_id = ANY($2) order by event_sequence", pq.Array(tt.res.eventsRes.aggTypes), pq.Array(tt.res.eventsRes.aggIDs))
|
||||||
|
if err != nil {
|
||||||
|
t.Error("unable to query inserted rows: ", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
var count int
|
||||||
|
|
||||||
|
for rows.Next() {
|
||||||
|
count++
|
||||||
|
data := make(Data, 0)
|
||||||
|
|
||||||
|
err := rows.Scan(&data)
|
||||||
|
if err != nil {
|
||||||
|
t.Error("unable to query inserted rows: ", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
t.Logf("inserted data: %v", string(data))
|
||||||
|
}
|
||||||
|
if count != tt.res.eventsRes.pushedEventsCount {
|
||||||
|
t.Errorf("expected push count %d got %d", tt.res.eventsRes.pushedEventsCount, count)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCRDB_query_events(t *testing.T) {
|
||||||
|
type args struct {
|
||||||
|
searchQuery *repository.SearchQuery
|
||||||
|
}
|
||||||
|
type fields struct {
|
||||||
|
existingEvents []*repository.Event
|
||||||
|
}
|
||||||
|
type res struct {
|
||||||
|
events []*repository.Event
|
||||||
|
}
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
fields fields
|
||||||
|
args args
|
||||||
|
res res
|
||||||
|
wantErr bool
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
name: "aggregate type filter no events",
|
||||||
|
args: args{
|
||||||
|
searchQuery: &repository.SearchQuery{},
|
||||||
|
},
|
||||||
|
fields: fields{
|
||||||
|
existingEvents: []*repository.Event{},
|
||||||
|
},
|
||||||
|
res: res{
|
||||||
|
events: []*repository.Event{},
|
||||||
|
},
|
||||||
|
wantErr: false,
|
||||||
|
},
|
||||||
|
// {
|
||||||
|
// name: "aggregate type filter events found",
|
||||||
|
// args: args{
|
||||||
|
// searchQuery: &repository.SearchQuery{},
|
||||||
|
// },
|
||||||
|
// fields: fields{
|
||||||
|
// existingEvents: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// res: res{
|
||||||
|
// events: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// wantErr: false,
|
||||||
|
// },
|
||||||
|
// {
|
||||||
|
// name: "aggregate type and id filter events found",
|
||||||
|
// args: args{
|
||||||
|
// searchQuery: &repository.SearchQuery{},
|
||||||
|
// },
|
||||||
|
// fields: fields{
|
||||||
|
// existingEvents: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// res: res{
|
||||||
|
// events: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// wantErr: false,
|
||||||
|
// },
|
||||||
|
// {
|
||||||
|
// name: "sequence filter events found",
|
||||||
|
// args: args{
|
||||||
|
// searchQuery: &repository.SearchQuery{},
|
||||||
|
// },
|
||||||
|
// fields: fields{
|
||||||
|
// existingEvents: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// res: res{
|
||||||
|
// events: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// wantErr: false,
|
||||||
|
// },
|
||||||
|
// {
|
||||||
|
// name: "resource owner filter events found",
|
||||||
|
// args: args{
|
||||||
|
// searchQuery: &repository.SearchQuery{},
|
||||||
|
// },
|
||||||
|
// fields: fields{
|
||||||
|
// existingEvents: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// res: res{
|
||||||
|
// events: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// wantErr: false,
|
||||||
|
// },
|
||||||
|
// {
|
||||||
|
// name: "editor service filter events found",
|
||||||
|
// args: args{
|
||||||
|
// searchQuery: &repository.SearchQuery{},
|
||||||
|
// },
|
||||||
|
// fields: fields{
|
||||||
|
// existingEvents: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// res: res{
|
||||||
|
// events: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// wantErr: false,
|
||||||
|
// },
|
||||||
|
// {
|
||||||
|
// name: "editor user filter events found",
|
||||||
|
// args: args{
|
||||||
|
// searchQuery: &repository.SearchQuery{},
|
||||||
|
// },
|
||||||
|
// fields: fields{
|
||||||
|
// existingEvents: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// res: res{
|
||||||
|
// events: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// wantErr: false,
|
||||||
|
// },
|
||||||
|
// {
|
||||||
|
// name: "event type filter events found",
|
||||||
|
// args: args{
|
||||||
|
// searchQuery: &repository.SearchQuery{},
|
||||||
|
// },
|
||||||
|
// fields: fields{
|
||||||
|
// existingEvents: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// res: res{
|
||||||
|
// events: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// wantErr: false,
|
||||||
|
// },
|
||||||
|
// {
|
||||||
|
// name: "no filter events found",
|
||||||
|
// args: args{
|
||||||
|
// searchQuery: &repository.SearchQuery{},
|
||||||
|
// },
|
||||||
|
// fields: fields{
|
||||||
|
// existingEvents: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// res: res{
|
||||||
|
// events: []*repository.Event{},
|
||||||
|
// },
|
||||||
|
// wantErr: false,
|
||||||
|
// },
|
||||||
|
}
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
db := &CRDB{
|
||||||
|
client: testCRDBClient,
|
||||||
|
}
|
||||||
|
|
||||||
|
// setup initial data for query
|
||||||
|
if err := db.Push(context.Background(), tt.fields.existingEvents...); err != nil {
|
||||||
|
t.Errorf("error in setup = %v", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
events := []*repository.Event{}
|
||||||
|
if err := db.query(tt.args.searchQuery, &events); (err != nil) != tt.wantErr {
|
||||||
|
t.Errorf("CRDB.query() error = %v, wantErr %v", err, tt.wantErr)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func canceledCtx() context.Context {
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
return ctx
|
||||||
|
}
|
||||||
|
|
||||||
|
func combineEventLists(lists ...[]*repository.Event) []*repository.Event {
|
||||||
|
combined := make([]*repository.Event, 0)
|
||||||
|
for _, list := range lists {
|
||||||
|
combined = append(combined, list...)
|
||||||
|
}
|
||||||
|
return combined
|
||||||
}
|
}
|
||||||
|
|
||||||
func linkEvents(events ...*repository.Event) []*repository.Event {
|
func linkEvents(events ...*repository.Event) []*repository.Event {
|
||||||
@ -563,6 +997,21 @@ func linkEvents(events ...*repository.Event) []*repository.Event {
|
|||||||
return events
|
return events
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func generateEventForAggregate(aggregateType repository.AggregateType, aggregateID string, checkPrevious bool, previousSeq uint64) *repository.Event {
|
||||||
|
return &repository.Event{
|
||||||
|
AggregateID: aggregateID,
|
||||||
|
AggregateType: aggregateType,
|
||||||
|
CheckPreviousSequence: checkPrevious,
|
||||||
|
EditorService: "svc",
|
||||||
|
EditorUser: "user",
|
||||||
|
PreviousEvent: nil,
|
||||||
|
PreviousSequence: previousSeq,
|
||||||
|
ResourceOwner: "ro",
|
||||||
|
Type: "test.created",
|
||||||
|
Version: "v1",
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func generateEvent(t *testing.T, aggregateID string, checkPrevious bool, previousSeq uint64) *repository.Event {
|
func generateEvent(t *testing.T, aggregateID string, checkPrevious bool, previousSeq uint64) *repository.Event {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
return &repository.Event{
|
return &repository.Event{
|
||||||
@ -578,3 +1027,20 @@ func generateEvent(t *testing.T, aggregateID string, checkPrevious bool, previou
|
|||||||
Version: "v1",
|
Version: "v1",
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func generateEventWithData(t *testing.T, aggregateID string, checkPrevious bool, previousSeq uint64, data []byte) *repository.Event {
|
||||||
|
t.Helper()
|
||||||
|
return &repository.Event{
|
||||||
|
AggregateID: aggregateID,
|
||||||
|
AggregateType: repository.AggregateType(t.Name()),
|
||||||
|
CheckPreviousSequence: checkPrevious,
|
||||||
|
EditorService: "svc",
|
||||||
|
EditorUser: "user",
|
||||||
|
PreviousEvent: nil,
|
||||||
|
PreviousSequence: previousSeq,
|
||||||
|
ResourceOwner: "ro",
|
||||||
|
Type: "test.created",
|
||||||
|
Version: "v1",
|
||||||
|
Data: data,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
@ -1,7 +1,6 @@
|
|||||||
package sql
|
package sql
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
|
||||||
"database/sql"
|
"database/sql"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
@ -11,10 +10,8 @@ import (
|
|||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/caos/logging"
|
"github.com/caos/logging"
|
||||||
"github.com/caos/zitadel/internal/eventstore/v2/repository"
|
|
||||||
"github.com/cockroachdb/cockroach-go/v2/testserver"
|
"github.com/cockroachdb/cockroach-go/v2/testserver"
|
||||||
)
|
)
|
||||||
|
|
||||||
@ -46,156 +43,156 @@ func TestMain(m *testing.M) {
|
|||||||
os.Exit(m.Run())
|
os.Exit(m.Run())
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestInsert(t *testing.T) {
|
// func TestInsert(t *testing.T) {
|
||||||
crdb := &CRDB{client: testCRDBClient}
|
// crdb := &CRDB{client: testCRDBClient}
|
||||||
e1 := &repository.Event{
|
// e1 := &repository.Event{
|
||||||
AggregateID: "agg.id",
|
// AggregateID: "agg.id",
|
||||||
AggregateType: "agg.type",
|
// AggregateType: "agg.type",
|
||||||
CheckPreviousSequence: true,
|
// CheckPreviousSequence: true,
|
||||||
EditorService: "edi.svc",
|
// EditorService: "edi.svc",
|
||||||
EditorUser: "edi",
|
// EditorUser: "edi",
|
||||||
ResourceOwner: "edit",
|
// ResourceOwner: "edit",
|
||||||
Type: "type",
|
// Type: "type",
|
||||||
Version: "v1",
|
// Version: "v1",
|
||||||
}
|
// }
|
||||||
events := []*repository.Event{
|
// events := []*repository.Event{
|
||||||
e1,
|
// e1,
|
||||||
{
|
// {
|
||||||
AggregateID: "agg.id",
|
// AggregateID: "agg.id",
|
||||||
AggregateType: "agg.type",
|
// AggregateType: "agg.type",
|
||||||
CheckPreviousSequence: true,
|
// CheckPreviousSequence: true,
|
||||||
EditorService: "edi.svc",
|
// EditorService: "edi.svc",
|
||||||
EditorUser: "edi",
|
// EditorUser: "edi",
|
||||||
ResourceOwner: "edit",
|
// ResourceOwner: "edit",
|
||||||
Type: "type",
|
// Type: "type",
|
||||||
Version: "v1",
|
// Version: "v1",
|
||||||
CreationDate: time.Now().Add(-2 * time.Second),
|
// CreationDate: time.Now().Add(-2 * time.Second),
|
||||||
PreviousEvent: e1,
|
// PreviousEvent: e1,
|
||||||
},
|
// },
|
||||||
{
|
// {
|
||||||
AggregateID: "agg.id",
|
// AggregateID: "agg.id",
|
||||||
AggregateType: "agg.type",
|
// AggregateType: "agg.type",
|
||||||
CheckPreviousSequence: false,
|
// CheckPreviousSequence: false,
|
||||||
EditorService: "edi.svc",
|
// EditorService: "edi.svc",
|
||||||
EditorUser: "edi",
|
// EditorUser: "edi",
|
||||||
ResourceOwner: "edit",
|
// ResourceOwner: "edit",
|
||||||
Type: "type",
|
// Type: "type",
|
||||||
Version: "v1",
|
// Version: "v1",
|
||||||
CreationDate: time.Now().Add(-500 * time.Millisecond),
|
// CreationDate: time.Now().Add(-500 * time.Millisecond),
|
||||||
},
|
// },
|
||||||
{
|
// {
|
||||||
AggregateID: "agg.id2",
|
// AggregateID: "agg.id2",
|
||||||
AggregateType: "agg.type",
|
// AggregateType: "agg.type",
|
||||||
CheckPreviousSequence: true,
|
// CheckPreviousSequence: true,
|
||||||
EditorService: "edi.svc",
|
// EditorService: "edi.svc",
|
||||||
EditorUser: "edi",
|
// EditorUser: "edi",
|
||||||
ResourceOwner: "edit",
|
// ResourceOwner: "edit",
|
||||||
Type: "type",
|
// Type: "type",
|
||||||
Version: "v1",
|
// Version: "v1",
|
||||||
CreationDate: time.Now().Add(-500 * time.Millisecond),
|
// CreationDate: time.Now().Add(-500 * time.Millisecond),
|
||||||
},
|
// },
|
||||||
{
|
// {
|
||||||
AggregateID: "agg.id3",
|
// AggregateID: "agg.id3",
|
||||||
AggregateType: "agg.type",
|
// AggregateType: "agg.type",
|
||||||
CheckPreviousSequence: false,
|
// CheckPreviousSequence: false,
|
||||||
EditorService: "edi.svc",
|
// EditorService: "edi.svc",
|
||||||
EditorUser: "edi",
|
// EditorUser: "edi",
|
||||||
ResourceOwner: "edit",
|
// ResourceOwner: "edit",
|
||||||
Type: "type",
|
// Type: "type",
|
||||||
Version: "v1",
|
// Version: "v1",
|
||||||
CreationDate: time.Now().Add(-500 * time.Millisecond),
|
// CreationDate: time.Now().Add(-500 * time.Millisecond),
|
||||||
},
|
// },
|
||||||
{
|
// {
|
||||||
AggregateID: "agg.id",
|
// AggregateID: "agg.id",
|
||||||
AggregateType: "agg.type",
|
// AggregateType: "agg.type",
|
||||||
CheckPreviousSequence: false,
|
// CheckPreviousSequence: false,
|
||||||
EditorService: "edi.svc",
|
// EditorService: "edi.svc",
|
||||||
EditorUser: "edi",
|
// EditorUser: "edi",
|
||||||
ResourceOwner: "edit",
|
// ResourceOwner: "edit",
|
||||||
Type: "type",
|
// Type: "type",
|
||||||
Version: "v1",
|
// Version: "v1",
|
||||||
CreationDate: time.Now().Add(-500 * time.Millisecond),
|
// CreationDate: time.Now().Add(-500 * time.Millisecond),
|
||||||
},
|
// },
|
||||||
// {
|
// // {
|
||||||
// AggregateID: "agg.id4",
|
// // AggregateID: "agg.id4",
|
||||||
// AggregateType: "agg.type",
|
// // AggregateType: "agg.type",
|
||||||
// CheckPreviousSequence: false,
|
// // CheckPreviousSequence: false,
|
||||||
// EditorService: "edi.svc",
|
// // EditorService: "edi.svc",
|
||||||
// EditorUser: "edi",
|
// // EditorUser: "edi",
|
||||||
// ResourceOwner: "edit",
|
// // ResourceOwner: "edit",
|
||||||
// Type: "type",
|
// // Type: "type",
|
||||||
// Version: "v1",
|
// // Version: "v1",
|
||||||
// CreationDate: time.Now().Add(-500 * time.Millisecond),
|
// // CreationDate: time.Now().Add(-500 * time.Millisecond),
|
||||||
// PreviousEvent: e1,
|
// // PreviousEvent: e1,
|
||||||
// },
|
// // },
|
||||||
//fail because wrong previous event
|
// //fail because wrong previous event
|
||||||
// {
|
// // {
|
||||||
// AggregateID: "agg.id2",
|
// // AggregateID: "agg.id2",
|
||||||
// AggregateType: "agg.type",
|
// // AggregateType: "agg.type",
|
||||||
// CheckPreviousSequence: true,
|
// // CheckPreviousSequence: true,
|
||||||
// EditorService: "edi.svc",
|
// // EditorService: "edi.svc",
|
||||||
// EditorUser: "edi",
|
// // EditorUser: "edi",
|
||||||
// ResourceOwner: "edit",
|
// // ResourceOwner: "edit",
|
||||||
// Type: "type",
|
// // Type: "type",
|
||||||
// Version: "v1",
|
// // Version: "v1",
|
||||||
// CreationDate: time.Now().Add(-500 * time.Millisecond),
|
// // CreationDate: time.Now().Add(-500 * time.Millisecond),
|
||||||
// PreviousEvent: e1,
|
// // PreviousEvent: e1,
|
||||||
// },
|
// // },
|
||||||
}
|
// }
|
||||||
fmt.Println("==============")
|
// fmt.Println("==============")
|
||||||
err := crdb.Push(context.Background(), events...)
|
// err := crdb.Push(context.Background(), events...)
|
||||||
if err != nil {
|
// if err != nil {
|
||||||
t.Error(err)
|
// t.Error(err)
|
||||||
return
|
// return
|
||||||
}
|
// }
|
||||||
|
|
||||||
fmt.Println("inserted ts:")
|
// fmt.Println("inserted ts:")
|
||||||
for _, event := range events {
|
// for _, event := range events {
|
||||||
fmt.Printf("%+v\n", event)
|
// fmt.Printf("%+v\n", event)
|
||||||
}
|
// }
|
||||||
|
|
||||||
fmt.Println("====================")
|
// fmt.Println("====================")
|
||||||
rows, err := testCRDBClient.Query("select * from eventstore.events order by event_sequence")
|
// rows, err := testCRDBClient.Query("select * from eventstore.events order by event_sequence")
|
||||||
defer rows.Close()
|
// defer rows.Close()
|
||||||
fmt.Println(err)
|
// fmt.Println(err)
|
||||||
|
|
||||||
for rows.Next() {
|
// for rows.Next() {
|
||||||
i := make([]interface{}, 12)
|
// i := make([]interface{}, 12)
|
||||||
var id string
|
// var id string
|
||||||
rows.Scan(&id, &i[1], &i[2], &i[3], &i[4], &i[5], &i[6], &i[7], &i[8], &i[9], &i[10], &i[11])
|
// rows.Scan(&id, &i[1], &i[2], &i[3], &i[4], &i[5], &i[6], &i[7], &i[8], &i[9], &i[10], &i[11])
|
||||||
i[0] = id
|
// i[0] = id
|
||||||
|
|
||||||
fmt.Println(i)
|
// fmt.Println(i)
|
||||||
}
|
// }
|
||||||
fmt.Println("====================")
|
// fmt.Println("====================")
|
||||||
filtededEvents, err := crdb.Filter(context.Background(), &repository.SearchQuery{
|
// filtededEvents, err := crdb.Filter(context.Background(), &repository.SearchQuery{
|
||||||
Columns: repository.ColumnsEvent,
|
// Columns: repository.ColumnsEvent,
|
||||||
Filters: []*repository.Filter{
|
// Filters: []*repository.Filter{
|
||||||
{
|
// {
|
||||||
Field: repository.FieldAggregateType,
|
// Field: repository.FieldAggregateType,
|
||||||
Operation: repository.OperationEquals,
|
// Operation: repository.OperationEquals,
|
||||||
Value: repository.AggregateType("agg.type"),
|
// Value: repository.AggregateType("agg.type"),
|
||||||
},
|
// },
|
||||||
},
|
// },
|
||||||
})
|
// })
|
||||||
fmt.Println(err)
|
// fmt.Println(err)
|
||||||
|
|
||||||
for _, event := range filtededEvents {
|
// for _, event := range filtededEvents {
|
||||||
fmt.Printf("%+v\n", event)
|
// fmt.Printf("%+v\n", event)
|
||||||
}
|
// }
|
||||||
fmt.Println("====================")
|
// fmt.Println("====================")
|
||||||
rows, err = testCRDBClient.Query("select max(event_sequence), count(*) from eventstore.events where aggregate_type = 'agg.type' and aggregate_id = 'agg.id'")
|
// rows, err = testCRDBClient.Query("select max(event_sequence), count(*) from eventstore.events where aggregate_type = 'agg.type' and aggregate_id = 'agg.id'")
|
||||||
defer rows.Close()
|
// defer rows.Close()
|
||||||
fmt.Println(err)
|
// fmt.Println(err)
|
||||||
|
|
||||||
for rows.Next() {
|
// for rows.Next() {
|
||||||
i := make([]interface{}, 2)
|
// i := make([]interface{}, 2)
|
||||||
rows.Scan(&i[0], &i[1])
|
// rows.Scan(&i[0], &i[1])
|
||||||
|
|
||||||
fmt.Println(i)
|
// fmt.Println(i)
|
||||||
}
|
// }
|
||||||
t.Fail()
|
// t.Fail()
|
||||||
}
|
// }
|
||||||
|
|
||||||
func executeMigrations() error {
|
func executeMigrations() error {
|
||||||
files, err := migrationFilePaths()
|
files, err := migrationFilePaths()
|
||||||
|
Loading…
x
Reference in New Issue
Block a user