mirror of
https://github.com/zitadel/zitadel.git
synced 2024-12-15 04:18:01 +00:00
b5564572bc
This implementation increases parallel write capabilities of the eventstore. Please have a look at the technical advisories: [05](https://zitadel.com/docs/support/advisory/a10005) and [06](https://zitadel.com/docs/support/advisory/a10006). The implementation of eventstore.push is rewritten and stored events are migrated to a new table `eventstore.events2`. If you are using cockroach: make sure that the database user of ZITADEL has `VIEWACTIVITY` grant. This is used to query events.
108 lines
2.6 KiB
Go
108 lines
2.6 KiB
Go
package eventstore
|
|
|
|
import (
|
|
"encoding/json"
|
|
"reflect"
|
|
"time"
|
|
|
|
"github.com/zitadel/zitadel/internal/errors"
|
|
)
|
|
|
|
type action interface {
|
|
Aggregate() *Aggregate
|
|
|
|
// Creator is the userid of the user which created the action
|
|
Creator() string
|
|
// Type describes the action
|
|
Type() EventType
|
|
// Revision of the action
|
|
Revision() uint16
|
|
}
|
|
|
|
// Command is the intend to store an event into the eventstore
|
|
type Command interface {
|
|
action
|
|
// Payload returns the payload of the event. It represent the changed fields by the event
|
|
// valid types are:
|
|
// * nil: no payload
|
|
// * struct: which can be marshalled to json
|
|
// * pointer: to struct which can be marshalled to json
|
|
// * []byte: json marshalled data
|
|
Payload() any
|
|
// UniqueConstraints should be added for unique attributes of an event, if nil constraints will not be checked
|
|
UniqueConstraints() []*UniqueConstraint
|
|
}
|
|
|
|
// Event is a stored activity
|
|
type Event interface {
|
|
action
|
|
|
|
// Sequence of the event in the aggregate
|
|
Sequence() uint64
|
|
// CreatedAt is the time the event was created at
|
|
CreatedAt() time.Time
|
|
// Position is the global position of the event
|
|
Position() float64
|
|
|
|
// Unmarshal parses the payload and stores the result
|
|
// in the value pointed to by ptr. If ptr is nil or not a pointer,
|
|
// Unmarshal returns an error
|
|
Unmarshal(ptr any) error
|
|
|
|
// Deprecated: only use for migration
|
|
DataAsBytes() []byte
|
|
}
|
|
|
|
type EventType string
|
|
|
|
func EventData(event Command) ([]byte, error) {
|
|
switch data := event.Payload().(type) {
|
|
case nil:
|
|
return nil, nil
|
|
case []byte:
|
|
if json.Valid(data) {
|
|
return data, nil
|
|
}
|
|
return nil, errors.ThrowInvalidArgument(nil, "V2-6SbbS", "data bytes are not json")
|
|
}
|
|
dataType := reflect.TypeOf(event.Payload())
|
|
if dataType.Kind() == reflect.Ptr {
|
|
dataType = dataType.Elem()
|
|
}
|
|
if dataType.Kind() == reflect.Struct {
|
|
dataBytes, err := json.Marshal(event.Payload())
|
|
if err != nil {
|
|
return nil, errors.ThrowInvalidArgument(err, "V2-xG87M", "could not marshal data")
|
|
}
|
|
return dataBytes, nil
|
|
}
|
|
return nil, errors.ThrowInvalidArgument(nil, "V2-91NRm", "wrong type of event data")
|
|
}
|
|
|
|
type BaseEventSetter[T any] interface {
|
|
Event
|
|
SetBaseEvent(*BaseEvent)
|
|
*T
|
|
}
|
|
|
|
func GenericEventMapper[T any, PT BaseEventSetter[T]](event Event) (Event, error) {
|
|
e := PT(new(T))
|
|
e.SetBaseEvent(BaseEventFromRepo(event))
|
|
|
|
err := event.Unmarshal(e)
|
|
if err != nil {
|
|
return nil, errors.ThrowInternal(err, "ES-Thai6", "unable to unmarshal event")
|
|
}
|
|
|
|
return e, nil
|
|
}
|
|
|
|
func isEventTypes(event Event, types ...EventType) bool {
|
|
for _, typ := range types {
|
|
if event.Type() == typ {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|