feat: add debug events API (#8533)

# Which Problems Are Solved

Add a debug API which allows pushing a set of events to be reduced in a
dedicated projection.
The events can carry a sleep duration which simulates a slow query
during projection handling.

# How the Problems Are Solved

- `CreateDebugEvents` allows pushing multiple events which simulate the
lifecycle of a resource. Each event has a `projectionSleep` field, which
issues a `pg_sleep()` statement query in the projection handler :
  - Add
  - Change
  - Remove
- `ListDebugEventsStates` list the current state of the projection,
optionally with a Trigger
- `GetDebugEventsStateByID` get the current state of the aggregate ID in
the projection, optionally with a Trigger


# Additional Changes

- none

# Additional Context

-  Allows reproduction of https://github.com/zitadel/zitadel/issues/8517
This commit is contained in:
Tim Möhlmann
2024-09-11 11:24:00 +03:00
committed by GitHub
parent a569501108
commit 3aba942162
21 changed files with 1404 additions and 0 deletions

View File

@@ -0,0 +1,82 @@
package command
import (
"context"
"time"
"github.com/zitadel/zitadel/internal/api/authz"
"github.com/zitadel/zitadel/internal/domain"
"github.com/zitadel/zitadel/internal/eventstore"
"github.com/zitadel/zitadel/internal/repository/debug_events"
"github.com/zitadel/zitadel/internal/telemetry/tracing"
"github.com/zitadel/zitadel/internal/zerrors"
)
type DebugEvents struct {
AggregateID string
Events []DebugEvent
}
type DebugEvent interface {
isADebugEvent()
}
type DebugEventAdded struct {
ProjectionSleep time.Duration
Blob *string
}
type DebugEventChanged struct {
ProjectionSleep time.Duration
Blob *string
}
type DebugEventRemoved struct {
ProjectionSleep time.Duration
}
func (DebugEventAdded) isADebugEvent() {}
func (DebugEventChanged) isADebugEvent() {}
func (DebugEventRemoved) isADebugEvent() {}
func (c *Commands) CreateDebugEvents(ctx context.Context, dbe *DebugEvents) (_ *domain.ObjectDetails, err error) {
ctx, span := tracing.NewSpan(ctx)
defer func() { span.EndWithError(err) }()
model := NewDebugEventsWriteModel(dbe.AggregateID, authz.GetInstance(ctx).InstanceID())
if err = c.eventstore.FilterToQueryReducer(ctx, model); err != nil {
return nil, err
}
aggr := debug_events.AggregateFromWriteModel(ctx, &model.WriteModel)
cmds := make([]eventstore.Command, len(dbe.Events))
for i, event := range dbe.Events {
var cmd eventstore.Command
switch e := event.(type) {
case DebugEventAdded:
if model.State.Exists() {
return nil, zerrors.ThrowAlreadyExists(nil, "COMMAND-Aex6j", "debug aggregate already exists")
}
cmd = debug_events.NewAddedEvent(ctx, aggr, e.ProjectionSleep, e.Blob)
case DebugEventChanged:
if !model.State.Exists() {
return nil, zerrors.ThrowNotFound(nil, "COMMAND-Thie6", "debug aggregate not found")
}
cmd = debug_events.NewChangedEvent(ctx, aggr, e.ProjectionSleep, e.Blob)
case DebugEventRemoved:
if !model.State.Exists() {
return nil, zerrors.ThrowNotFound(nil, "COMMAND-Ohna9", "debug aggregate not found")
}
cmd = debug_events.NewRemovedEvent(ctx, aggr, e.ProjectionSleep)
}
cmds[i] = cmd
// be sure the state of the last event is reduced before handling the next one.
model.reduceEvent(cmd.(eventstore.Event))
}
events, err := c.eventstore.Push(ctx, cmds...)
if err != nil {
return nil, err
}
return pushedEventsToObjectDetails(events), nil
}

View File

@@ -0,0 +1,68 @@
package command
import (
"github.com/zitadel/zitadel/internal/domain"
"github.com/zitadel/zitadel/internal/eventstore"
debug "github.com/zitadel/zitadel/internal/repository/debug_events"
)
type DebugEventsWriteModel struct {
eventstore.WriteModel
State domain.DebugEventsState
Blob string
}
func NewDebugEventsWriteModel(aggregateID, resourceOwner string) *DebugEventsWriteModel {
return &DebugEventsWriteModel{
WriteModel: eventstore.WriteModel{
AggregateID: aggregateID,
ResourceOwner: resourceOwner,
},
}
}
func (wm *DebugEventsWriteModel) AppendEvents(events ...eventstore.Event) {
wm.WriteModel.AppendEvents(events...)
}
func (wm *DebugEventsWriteModel) Reduce() error {
for _, event := range wm.Events {
wm.reduceEvent(event)
}
return wm.WriteModel.Reduce()
}
func (wm *DebugEventsWriteModel) reduceEvent(event eventstore.Event) {
if event.Aggregate().ID != wm.AggregateID {
return
}
switch e := event.(type) {
case *debug.AddedEvent:
wm.State = domain.DebugEventsStateInitial
if e.Blob != nil {
wm.Blob = *e.Blob
}
case *debug.ChangedEvent:
wm.State = domain.DebugEventsStateChanged
if e.Blob != nil {
wm.Blob = *e.Blob
}
case *debug.RemovedEvent:
wm.State = domain.DebugEventsStateRemoved
wm.Blob = ""
}
}
func (wm *DebugEventsWriteModel) Query() *eventstore.SearchQueryBuilder {
return eventstore.NewSearchQueryBuilder(eventstore.ColumnsEvent).
ResourceOwner(wm.ResourceOwner).
AddQuery().
AggregateTypes(debug.AggregateType).
AggregateIDs(wm.AggregateID).
EventTypes(
debug.AddedEventType,
debug.ChangedEventType,
debug.RemovedEventType,
).
Builder()
}

View File

@@ -0,0 +1,340 @@
package command
import (
"io"
"testing"
"time"
"github.com/muhlemmer/gu"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/zitadel/zitadel/internal/api/authz"
"github.com/zitadel/zitadel/internal/domain"
"github.com/zitadel/zitadel/internal/eventstore"
"github.com/zitadel/zitadel/internal/repository/debug_events"
"github.com/zitadel/zitadel/internal/zerrors"
)
func TestCommands_CreateDebugEvents(t *testing.T) {
ctx := authz.NewMockContextWithPermissions("instance1", "org1", "user1", nil)
type fields struct {
eventstore func(*testing.T) *eventstore.Eventstore
}
type args struct {
dbe *DebugEvents
}
tests := []struct {
name string
fields fields
args args
want *domain.ObjectDetails
wantErr error
}{
{
name: "filter error",
fields: fields{
eventstore: expectEventstore(
expectFilterError(io.ErrClosedPipe),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventAdded{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("a"),
},
},
}},
wantErr: io.ErrClosedPipe,
},
{
name: "already exists",
fields: fields{
eventstore: expectEventstore(
expectFilter(
eventFromEventPusher(
debug_events.NewAddedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("a"),
),
),
),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventAdded{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("a"),
},
},
}},
wantErr: zerrors.ThrowAlreadyExists(nil, "COMMAND-Aex6j", "debug aggregate already exists"),
},
{
name: "double added event, already exists",
fields: fields{
eventstore: expectEventstore(
expectFilter(),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventAdded{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("a"),
},
DebugEventAdded{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("a"),
},
},
}},
wantErr: zerrors.ThrowAlreadyExists(nil, "COMMAND-Aex6j", "debug aggregate already exists"),
},
{
name: "changed event, not found",
fields: fields{
eventstore: expectEventstore(
expectFilter(),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventChanged{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("a"),
},
},
}},
wantErr: zerrors.ThrowNotFound(nil, "COMMAND-Thie6", "debug aggregate not found"),
},
{
name: "removed event, not found",
fields: fields{
eventstore: expectEventstore(
expectFilter(),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventRemoved{
ProjectionSleep: time.Millisecond,
},
},
}},
wantErr: zerrors.ThrowNotFound(nil, "COMMAND-Ohna9", "debug aggregate not found"),
},
{
name: "changed after removed event, not found",
fields: fields{
eventstore: expectEventstore(
expectFilter(
eventFromEventPusher(
debug_events.NewAddedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("a"),
),
),
),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventRemoved{
ProjectionSleep: time.Millisecond,
},
DebugEventChanged{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("a"),
},
},
}},
wantErr: zerrors.ThrowNotFound(nil, "COMMAND-Thie6", "debug aggregate not found"),
},
{
name: "double removed event, not found",
fields: fields{
eventstore: expectEventstore(
expectFilter(
eventFromEventPusher(
debug_events.NewAddedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("a"),
),
),
),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventRemoved{
ProjectionSleep: time.Millisecond,
},
DebugEventRemoved{
ProjectionSleep: time.Millisecond,
},
},
}},
wantErr: zerrors.ThrowNotFound(nil, "COMMAND-Ohna9", "debug aggregate not found"),
},
{
name: "added, ok",
fields: fields{
eventstore: expectEventstore(
expectFilter(),
expectPush(
debug_events.NewAddedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("a"),
),
),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventAdded{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("a"),
},
},
}},
want: &domain.ObjectDetails{
ResourceOwner: "instance1",
},
},
{
name: "changed, ok",
fields: fields{
eventstore: expectEventstore(
expectFilter(
eventFromEventPusher(
debug_events.NewAddedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("a"),
),
),
),
expectPush(
debug_events.NewChangedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("b"),
),
),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventChanged{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("b"),
},
},
}},
want: &domain.ObjectDetails{
ResourceOwner: "instance1",
},
},
{
name: "removed, ok",
fields: fields{
eventstore: expectEventstore(
expectFilter(
eventFromEventPusher(
debug_events.NewAddedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("a"),
),
),
),
expectPush(
debug_events.NewRemovedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond,
),
),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventRemoved{
ProjectionSleep: time.Millisecond,
},
},
}},
want: &domain.ObjectDetails{
ResourceOwner: "instance1",
},
},
{
name: "added, changed, changed, removed ok",
fields: fields{
eventstore: expectEventstore(
expectFilter(),
expectPush(
debug_events.NewAddedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("a"),
),
debug_events.NewChangedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("b"),
),
debug_events.NewChangedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond, gu.Ptr("c"),
),
debug_events.NewRemovedEvent(
ctx, debug_events.NewAggregate("dbg1", "instance1"),
time.Millisecond,
),
),
),
},
args: args{&DebugEvents{
AggregateID: "dbg1",
Events: []DebugEvent{
DebugEventAdded{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("a"),
},
DebugEventChanged{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("b"),
},
DebugEventChanged{
ProjectionSleep: time.Millisecond,
Blob: gu.Ptr("c"),
},
DebugEventRemoved{
ProjectionSleep: time.Millisecond,
},
},
}},
want: &domain.ObjectDetails{
ResourceOwner: "instance1",
},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
c := &Commands{
eventstore: tt.fields.eventstore(t),
}
got, err := c.CreateDebugEvents(ctx, tt.args.dbe)
require.ErrorIs(t, err, tt.wantErr)
assert.Equal(t, tt.want, got)
})
}
}