2023-10-19 12:19:10 +02:00
|
|
|
package eventstore_test
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"testing"
|
|
|
|
|
|
|
|
"github.com/zitadel/zitadel/internal/eventstore"
|
|
|
|
)
|
|
|
|
|
|
|
|
func TestCRDB_Filter(t *testing.T) {
|
|
|
|
type args struct {
|
|
|
|
searchQuery *eventstore.SearchQueryBuilder
|
|
|
|
}
|
|
|
|
type fields struct {
|
|
|
|
existingEvents []eventstore.Command
|
|
|
|
}
|
|
|
|
type res struct {
|
|
|
|
eventCount int
|
|
|
|
}
|
|
|
|
tests := []struct {
|
|
|
|
name string
|
|
|
|
fields fields
|
|
|
|
args args
|
|
|
|
res res
|
|
|
|
wantErr bool
|
|
|
|
}{
|
|
|
|
{
|
|
|
|
name: "aggregate type filter no events",
|
|
|
|
args: args{
|
|
|
|
searchQuery: eventstore.NewSearchQueryBuilder(eventstore.ColumnsEvent).
|
|
|
|
AddQuery().
|
|
|
|
AggregateTypes("not found").
|
|
|
|
Builder(),
|
|
|
|
},
|
|
|
|
fields: fields{
|
|
|
|
existingEvents: []eventstore.Command{
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "300"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "300"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "300"),
|
|
|
|
},
|
|
|
|
},
|
|
|
|
res: res{
|
|
|
|
eventCount: 0,
|
|
|
|
},
|
|
|
|
wantErr: false,
|
|
|
|
},
|
|
|
|
{
|
|
|
|
name: "aggregate type and id filter events found",
|
|
|
|
args: args{
|
|
|
|
searchQuery: eventstore.NewSearchQueryBuilder(eventstore.ColumnsEvent).
|
|
|
|
AddQuery().
|
|
|
|
AggregateTypes(eventstore.AggregateType(t.Name())).
|
|
|
|
AggregateIDs("303").
|
|
|
|
Builder(),
|
|
|
|
},
|
|
|
|
fields: fields{
|
|
|
|
existingEvents: []eventstore.Command{
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "303"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "303"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "303"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "305"),
|
|
|
|
},
|
|
|
|
},
|
|
|
|
res: res{
|
|
|
|
eventCount: 3,
|
|
|
|
},
|
|
|
|
wantErr: false,
|
|
|
|
},
|
feat(eventstore): exclude aggregate IDs when event_type occurred (#8940)
# Which Problems Are Solved
For truly event-based notification handler, we need to be able to filter
out events of aggregates which are already handled. For example when an
event like `notify.success` or `notify.failed` was created on an
aggregate, we no longer require events from that aggregate ID.
# How the Problems Are Solved
Extend the query builder to use a `NOT IN` clause which excludes
aggregate IDs when they have certain events for a certain aggregate
type. For optimization and proper index usages, certain filters are
inherited from the parent query, such as:
- Instance ID
- Instance IDs
- Position offset
This is a prettified query as used by the unit tests:
```sql
SELECT created_at, event_type, "sequence", "position", payload, creator, "owner", instance_id, aggregate_type, aggregate_id, revision
FROM eventstore.events2
WHERE instance_id = $1
AND aggregate_type = $2
AND event_type = $3
AND "position" > $4
AND aggregate_id NOT IN (
SELECT aggregate_id
FROM eventstore.events2
WHERE aggregate_type = $5
AND event_type = ANY($6)
AND instance_id = $7
AND "position" > $8
)
ORDER BY "position" DESC, in_tx_order DESC
LIMIT $9
```
I used this query to run it against the `oidc_session` aggregate looking
for added events, excluding aggregates where a token was revoked,
against a recent position. It fully used index scans:
<details>
```json
[
{
"Plan": {
"Node Type": "Index Scan",
"Parallel Aware": false,
"Async Capable": false,
"Scan Direction": "Forward",
"Index Name": "es_projection",
"Relation Name": "events2",
"Alias": "events2",
"Actual Rows": 2,
"Actual Loops": 1,
"Index Cond": "((instance_id = '286399006995644420'::text) AND (aggregate_type = 'oidc_session'::text) AND (event_type = 'oidc_session.added'::text) AND (\"position\" > 1731582100.784168))",
"Rows Removed by Index Recheck": 0,
"Filter": "(NOT (hashed SubPlan 1))",
"Rows Removed by Filter": 1,
"Plans": [
{
"Node Type": "Index Scan",
"Parent Relationship": "SubPlan",
"Subplan Name": "SubPlan 1",
"Parallel Aware": false,
"Async Capable": false,
"Scan Direction": "Forward",
"Index Name": "es_projection",
"Relation Name": "events2",
"Alias": "events2_1",
"Actual Rows": 1,
"Actual Loops": 1,
"Index Cond": "((instance_id = '286399006995644420'::text) AND (aggregate_type = 'oidc_session'::text) AND (event_type = 'oidc_session.access_token.revoked'::text) AND (\"position\" > 1731582100.784168))",
"Rows Removed by Index Recheck": 0
}
]
},
"Triggers": [
]
}
]
```
</details>
# Additional Changes
- None
# Additional Context
- Related to https://github.com/zitadel/zitadel/issues/8931
---------
Co-authored-by: adlerhurst <silvan.reusser@gmail.com>
2024-11-25 17:25:11 +02:00
|
|
|
{
|
|
|
|
name: "exclude aggregate type and event type",
|
|
|
|
args: args{
|
|
|
|
searchQuery: eventstore.NewSearchQueryBuilder(eventstore.ColumnsEvent).
|
|
|
|
AddQuery().
|
|
|
|
AggregateTypes(eventstore.AggregateType(t.Name())).
|
|
|
|
Builder().
|
|
|
|
ExcludeAggregateIDs().
|
|
|
|
EventTypes("test.updated").
|
|
|
|
AggregateTypes(eventstore.AggregateType(t.Name())).
|
|
|
|
Builder(),
|
|
|
|
},
|
|
|
|
fields: fields{
|
|
|
|
existingEvents: []eventstore.Command{
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "306"),
|
|
|
|
generateCommand(
|
|
|
|
eventstore.AggregateType(t.Name()),
|
|
|
|
"306",
|
|
|
|
func(te *testEvent) {
|
|
|
|
te.EventType = "test.updated"
|
|
|
|
},
|
|
|
|
),
|
|
|
|
generateCommand(
|
|
|
|
eventstore.AggregateType(t.Name()),
|
|
|
|
"308",
|
|
|
|
),
|
|
|
|
},
|
|
|
|
},
|
|
|
|
res: res{
|
|
|
|
eventCount: 1,
|
|
|
|
},
|
|
|
|
wantErr: false,
|
|
|
|
},
|
2023-10-19 12:19:10 +02:00
|
|
|
}
|
|
|
|
for _, tt := range tests {
|
|
|
|
for querierName, querier := range queriers {
|
|
|
|
t.Run(querierName+"/"+tt.name, func(t *testing.T) {
|
|
|
|
t.Cleanup(cleanupEventstore(clients[querierName]))
|
|
|
|
|
|
|
|
db := eventstore.NewEventstore(
|
|
|
|
&eventstore.Config{
|
|
|
|
Querier: querier,
|
|
|
|
Pusher: pushers["v3(inmemory)"],
|
|
|
|
},
|
|
|
|
)
|
|
|
|
|
|
|
|
// 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, err := db.Filter(context.Background(), tt.args.searchQuery)
|
|
|
|
if (err != nil) != tt.wantErr {
|
|
|
|
t.Errorf("CRDB.query() error = %v, wantErr %v", err, tt.wantErr)
|
|
|
|
}
|
|
|
|
|
|
|
|
if len(events) != tt.res.eventCount {
|
|
|
|
t.Errorf("CRDB.query() expected event count: %d got %d", tt.res.eventCount, len(events))
|
|
|
|
}
|
|
|
|
})
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-09-24 19:43:29 +03:00
|
|
|
func TestCRDB_LatestSequence(t *testing.T) {
|
2023-10-19 12:19:10 +02:00
|
|
|
type args struct {
|
|
|
|
searchQuery *eventstore.SearchQueryBuilder
|
|
|
|
}
|
|
|
|
type fields struct {
|
|
|
|
existingEvents []eventstore.Command
|
|
|
|
}
|
|
|
|
type res struct {
|
2024-09-24 19:43:29 +03:00
|
|
|
sequence float64
|
2023-10-19 12:19:10 +02:00
|
|
|
}
|
|
|
|
tests := []struct {
|
|
|
|
name string
|
|
|
|
fields fields
|
|
|
|
args args
|
|
|
|
res res
|
|
|
|
wantErr bool
|
|
|
|
}{
|
|
|
|
{
|
|
|
|
name: "aggregate type filter no sequence",
|
|
|
|
args: args{
|
2024-09-24 19:43:29 +03:00
|
|
|
searchQuery: eventstore.NewSearchQueryBuilder(eventstore.ColumnsMaxSequence).
|
2023-10-19 12:19:10 +02:00
|
|
|
AddQuery().
|
|
|
|
AggregateTypes("not found").
|
|
|
|
Builder(),
|
|
|
|
},
|
|
|
|
fields: fields{
|
|
|
|
existingEvents: []eventstore.Command{
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "400"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "400"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "400"),
|
|
|
|
},
|
|
|
|
},
|
|
|
|
wantErr: false,
|
|
|
|
},
|
|
|
|
{
|
|
|
|
name: "aggregate type filter sequence",
|
|
|
|
args: args{
|
2024-09-24 19:43:29 +03:00
|
|
|
searchQuery: eventstore.NewSearchQueryBuilder(eventstore.ColumnsMaxSequence).
|
2023-10-19 12:19:10 +02:00
|
|
|
AddQuery().
|
|
|
|
AggregateTypes(eventstore.AggregateType(t.Name())).
|
|
|
|
Builder(),
|
|
|
|
},
|
|
|
|
fields: fields{
|
|
|
|
existingEvents: []eventstore.Command{
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "401"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "401"),
|
|
|
|
generateCommand(eventstore.AggregateType(t.Name()), "401"),
|
|
|
|
},
|
|
|
|
},
|
|
|
|
wantErr: false,
|
|
|
|
},
|
|
|
|
}
|
|
|
|
for _, tt := range tests {
|
|
|
|
for querierName, querier := range queriers {
|
|
|
|
t.Run(querierName+"/"+tt.name, func(t *testing.T) {
|
|
|
|
t.Cleanup(cleanupEventstore(clients[querierName]))
|
|
|
|
|
|
|
|
db := eventstore.NewEventstore(
|
|
|
|
&eventstore.Config{
|
|
|
|
Querier: querier,
|
|
|
|
Pusher: pushers["v3(inmemory)"],
|
|
|
|
},
|
|
|
|
)
|
|
|
|
|
|
|
|
// setup initial data for query
|
|
|
|
_, err := db.Push(context.Background(), tt.fields.existingEvents...)
|
|
|
|
if err != nil {
|
|
|
|
t.Errorf("error in setup = %v", err)
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2024-09-24 19:43:29 +03:00
|
|
|
sequence, err := db.LatestSequence(context.Background(), tt.args.searchQuery)
|
2023-10-19 12:19:10 +02:00
|
|
|
if (err != nil) != tt.wantErr {
|
|
|
|
t.Errorf("CRDB.query() error = %v, wantErr %v", err, tt.wantErr)
|
|
|
|
}
|
2024-09-24 19:43:29 +03:00
|
|
|
if tt.res.sequence > sequence {
|
|
|
|
t.Errorf("CRDB.query() expected sequence: %v got %v", tt.res.sequence, sequence)
|
2023-10-19 12:19:10 +02:00
|
|
|
}
|
|
|
|
})
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|