diff --git a/server/app/boot_levels.go b/server/app/boot_levels.go index 172e27a0b..5b5bdca95 100644 --- a/server/app/boot_levels.go +++ b/server/app/boot_levels.go @@ -27,6 +27,7 @@ import ( "github.com/cortezaproject/corteza/server/pkg/eventbus" "github.com/cortezaproject/corteza/server/pkg/healthcheck" "github.com/cortezaproject/corteza/server/pkg/http" + "github.com/cortezaproject/corteza/server/pkg/id" "github.com/cortezaproject/corteza/server/pkg/locale" "github.com/cortezaproject/corteza/server/pkg/logger" "github.com/cortezaproject/corteza/server/pkg/mail" @@ -293,6 +294,9 @@ func (app *CortezaApp) InitServices(ctx context.Context) (err error) { return nil } + // @todo place this somewhere better + id.Init(ctx) + err = app.initEnvoy(ctx, app.Log) if err != nil { return diff --git a/server/pkg/id/sonyflake.go b/server/pkg/id/sonyflake.go index 1562e6fa6..b489d140c 100644 --- a/server/pkg/id/sonyflake.go +++ b/server/pkg/id/sonyflake.go @@ -1,11 +1,14 @@ package id import ( + "context" "fmt" - "github.com/sony/sonyflake" "os" "strconv" + "sync" "time" + + "github.com/sony/sonyflake" ) const ( @@ -33,54 +36,91 @@ const ( // envKeySimpleIncrement = "SONYFLAKE_SIMPLE_INCREMENT" ) -var sf *sonyflake.Sonyflake +var ( + initialized bool + idQueue = make(chan uint64, 512) + makerCount = 20 +) -func init() { - if err := initSonyflake(); err != nil { - panic(fmt.Errorf("sonyflake init failed: %w", err)) +func Init(ctx context.Context) { + if initialized { + return } + initialized = true + + wg := &sync.WaitGroup{} + wg.Add(makerCount) + + for i := 0; i < makerCount; i++ { + makeIDer(ctx, wg, idQueue, uint64(i+1)) + } + + wg.Wait() } -func initSonyflake() error { +func makeIDer(ctx context.Context, wg *sync.WaitGroup, qq chan uint64, thr uint64) { + go func() { + ss, err := initSonyflake(thr) + if err != nil { + panic(fmt.Errorf("sonyflake init failed: %w", err)) + } + + wg.Done() + for { + select { + case <-ctx.Done(): + return + + default: + id, err := ss.NextID() + if err != nil { + panic(err) + } + qq <- id + } + } + }() +} + +func initSonyflake(thr uint64) (out *sonyflake.Sonyflake, err error) { settings := sonyflake.Settings{ StartTime: time.Unix(jan1st2017, 0), } // check if SONYFLAKE_MACHINE_ID is set and create // MachineID function that returns it + var id uint64 if machineID := os.Getenv(envKeyMachineID); machineID != "" { // check if machineID is a valid uint16 - if id, err := strconv.ParseUint(machineID, 10, 16); err != nil { + if id, err = strconv.ParseUint(machineID, 10, 16); err != nil { // should crash right away - return fmt.Errorf( + err = fmt.Errorf( "could not use %s (%s), expecting uint16 (0..65535): %w", envKeyMachineID, machineID, err, ) + return } else { settings.MachineID = func() (uint16, error) { - return uint16(id), nil + return uint16(id + thr), nil } } + } else if thr > 0 { + settings.MachineID = func() (uint16, error) { + return uint16(thr), nil + } } - sf = sonyflake.NewSonyflake(settings) - _, err := sf.NextID() - return err -} + out = sonyflake.NewSonyflake(settings) + _, err = out.NextID() -// NextID returns uint64 ID, or panics -// -// See https://github.com/sony/sonyflake for details -func nextSonyflake() uint64 { - if id, err := sf.NextID(); err != nil { - panic(err) - } else { - return id - } + return } func Next() uint64 { - return nextSonyflake() + if !initialized { + panic("ID generator not initialized: call pkg/id.Init") + } + return <-idQueue } diff --git a/server/pkg/id/sonyflake_test.go b/server/pkg/id/sonyflake_test.go new file mode 100644 index 000000000..f327864a8 --- /dev/null +++ b/server/pkg/id/sonyflake_test.go @@ -0,0 +1,16 @@ +package id + +import ( + "testing" + + "github.com/cortezaproject/corteza/server/pkg/cli" +) + +func BenchmarkGenerator(b *testing.B) { + ctx := cli.Context() + Init(ctx) + + for n := 0; n < b.N; n++ { + Next() + } +}