mirror of
https://github.com/zitadel/zitadel.git
synced 2025-01-05 14:37:45 +00:00
dab5d9e756
# Which Problems Are Solved If many events are written to the same aggregate id it can happen that zitadel [starts to retry the push transaction](48ffc902cc/internal/eventstore/eventstore.go (L101)
) because [the locking behaviour](48ffc902cc/internal/eventstore/v3/sequence.go (L25)
) during push does compute the wrong sequence because newly committed events are not visible to the transaction. These events impact the current sequence. In cases with high command traffic on a single aggregate id this can have severe impact on general performance of zitadel. Because many connections of the `eventstore pusher` database pool are blocked by each other. # How the Problems Are Solved To improve the performance this locking mechanism was removed and the business logic of push is moved to sql functions which reduce network traffic and can be analyzed by the database before the actual push. For clients of the eventstore framework nothing changed. # Additional Changes - after a connection is established prefetches the newly added database types - `eventstore.BaseEvent` now returns the correct revision of the event # Additional Context - part of https://github.com/zitadel/zitadel/issues/8931 --------- Co-authored-by: Tim Möhlmann <tim+github@zitadel.com> Co-authored-by: Livio Spring <livio.a@gmail.com> Co-authored-by: Max Peintner <max@caos.ch> Co-authored-by: Elio Bischof <elio@zitadel.com> Co-authored-by: Stefan Benz <46600784+stebenz@users.noreply.github.com> Co-authored-by: Miguel Cabrerizo <30386061+doncicuto@users.noreply.github.com> Co-authored-by: Joakim Lodén <Loddan@users.noreply.github.com> Co-authored-by: Yxnt <Yxnt@users.noreply.github.com> Co-authored-by: Stefan Benz <stefan@caos.ch> Co-authored-by: Harsha Reddy <harsha.reddy@klaviyo.com> Co-authored-by: Zach H <zhirschtritt@gmail.com>
130 lines
3.9 KiB
Go
130 lines
3.9 KiB
Go
package dialect
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"reflect"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgtype"
|
|
)
|
|
|
|
var (
|
|
ErrNegativeRatio = errors.New("ratio cannot be negative")
|
|
ErrHighSumRatio = errors.New("sum of pusher and projection ratios must be < 1")
|
|
ErrIllegalMaxOpenConns = errors.New("MaxOpenConns of the database must be higher than 3 or 0 for unlimited")
|
|
ErrIllegalMaxIdleConns = errors.New("MaxIdleConns of the database must be higher than 3 or 0 for unlimited")
|
|
ErrInvalidPurpose = errors.New("DBPurpose out of range")
|
|
)
|
|
|
|
// ConnectionConfig defines the Max Open and Idle connections for a DB connection pool.
|
|
type ConnectionConfig struct {
|
|
MaxOpenConns,
|
|
MaxIdleConns uint32
|
|
AfterConnect []func(ctx context.Context, c *pgx.Conn) error
|
|
}
|
|
|
|
// takeRatio of MaxOpenConns and MaxIdleConns from config and returns
|
|
// a new ConnectionConfig with the resulting values.
|
|
func (c *ConnectionConfig) takeRatio(ratio float64) (*ConnectionConfig, error) {
|
|
if ratio < 0 {
|
|
return nil, ErrNegativeRatio
|
|
}
|
|
|
|
out := &ConnectionConfig{
|
|
MaxOpenConns: uint32(ratio * float64(c.MaxOpenConns)),
|
|
MaxIdleConns: uint32(ratio * float64(c.MaxIdleConns)),
|
|
AfterConnect: c.AfterConnect,
|
|
}
|
|
if c.MaxOpenConns != 0 && out.MaxOpenConns < 1 && ratio > 0 {
|
|
out.MaxOpenConns = 1
|
|
}
|
|
if c.MaxIdleConns != 0 && out.MaxIdleConns < 1 && ratio > 0 {
|
|
out.MaxIdleConns = 1
|
|
}
|
|
|
|
return out, nil
|
|
}
|
|
|
|
var afterConnectFuncs []func(ctx context.Context, c *pgx.Conn) error
|
|
|
|
func RegisterAfterConnect(f func(ctx context.Context, c *pgx.Conn) error) {
|
|
afterConnectFuncs = append(afterConnectFuncs, f)
|
|
}
|
|
|
|
func RegisterDefaultPgTypeVariants[T any](m *pgtype.Map, name, arrayName string) {
|
|
// T
|
|
var value T
|
|
m.RegisterDefaultPgType(value, name)
|
|
|
|
// *T
|
|
valueType := reflect.TypeOf(value)
|
|
m.RegisterDefaultPgType(reflect.New(valueType).Interface(), name)
|
|
|
|
// []T
|
|
sliceType := reflect.SliceOf(valueType)
|
|
m.RegisterDefaultPgType(reflect.MakeSlice(sliceType, 0, 0).Interface(), arrayName)
|
|
|
|
// *[]T
|
|
m.RegisterDefaultPgType(reflect.New(sliceType).Interface(), arrayName)
|
|
|
|
// []*T
|
|
sliceOfPointerType := reflect.SliceOf(reflect.TypeOf(reflect.New(valueType).Interface()))
|
|
m.RegisterDefaultPgType(reflect.MakeSlice(sliceOfPointerType, 0, 0).Interface(), arrayName)
|
|
|
|
// *[]*T
|
|
m.RegisterDefaultPgType(reflect.New(sliceOfPointerType).Interface(), arrayName)
|
|
}
|
|
|
|
// NewConnectionConfig calculates [ConnectionConfig] values from the passed ratios
|
|
// and returns the config applicable for the requested purpose.
|
|
//
|
|
// openConns and idleConns must be at least 3 or 0, which means no limit.
|
|
// The pusherRatio and spoolerRatio must be between 0 and 1.
|
|
func NewConnectionConfig(openConns, idleConns uint32, pusherRatio, projectionRatio float64, purpose DBPurpose) (*ConnectionConfig, error) {
|
|
if openConns != 0 && openConns < 3 {
|
|
return nil, ErrIllegalMaxOpenConns
|
|
}
|
|
if idleConns != 0 && idleConns < 3 {
|
|
return nil, ErrIllegalMaxIdleConns
|
|
}
|
|
if pusherRatio+projectionRatio >= 1 {
|
|
return nil, ErrHighSumRatio
|
|
}
|
|
|
|
queryConfig := &ConnectionConfig{
|
|
MaxOpenConns: openConns,
|
|
MaxIdleConns: idleConns,
|
|
AfterConnect: afterConnectFuncs,
|
|
}
|
|
pusherConfig, err := queryConfig.takeRatio(pusherRatio)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("event pusher: %w", err)
|
|
}
|
|
|
|
spoolerConfig, err := queryConfig.takeRatio(projectionRatio)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("projection spooler: %w", err)
|
|
}
|
|
|
|
// subtract the claimed amount
|
|
if queryConfig.MaxOpenConns > 0 {
|
|
queryConfig.MaxOpenConns -= pusherConfig.MaxOpenConns + spoolerConfig.MaxOpenConns
|
|
}
|
|
if queryConfig.MaxIdleConns > 0 {
|
|
queryConfig.MaxIdleConns -= pusherConfig.MaxIdleConns + spoolerConfig.MaxIdleConns
|
|
}
|
|
|
|
switch purpose {
|
|
case DBPurposeQuery:
|
|
return queryConfig, nil
|
|
case DBPurposeEventPusher:
|
|
return pusherConfig, nil
|
|
case DBPurposeProjectionSpooler:
|
|
return spoolerConfig, nil
|
|
default:
|
|
return nil, fmt.Errorf("%w: %v", ErrInvalidPurpose, purpose)
|
|
}
|
|
}
|