u
This commit is contained in:
1 parent
4ef00e1d5d
commit
9e00adf2fe
27 files changed
+1104
No files matched your search
@@ -0,0 +1,38 @@
|
||||
package bootstrap
|
||||
|
||||
import (
|
||||
"allapp-go-v2/internal/config"
|
||||
"allapp-go-v2/internal/pkg/logger"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
type App struct {
|
||||
Config *config.Config
|
||||
Db *pgxpool.Pool
|
||||
}
|
||||
|
||||
func NewApp() (*App, error) {
|
||||
cfg, err := config.Load()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 初始化日志
|
||||
if err := logger.Init(cfg.App.Env); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
InitUniqueId(cfg)
|
||||
|
||||
db, err := InitDB(cfg)
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &App{
|
||||
Config: cfg,
|
||||
Db: db,
|
||||
}, nil
|
||||
}
|
||||
@@ -0,0 +1,53 @@
|
||||
package bootstrap
|
||||
|
||||
import (
|
||||
"allapp-go-v2/internal/config"
|
||||
"allapp-go-v2/internal/pkg/logger"
|
||||
"context"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
func InitDB(cfg *config.Config) (*pgxpool.Pool, error) {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
pg := cfg.Postgres
|
||||
|
||||
// ssl mode
|
||||
sslMode := "disable"
|
||||
if pg.SslMode {
|
||||
sslMode = "require"
|
||||
}
|
||||
|
||||
// 转义密码(防止特殊字符问题)
|
||||
password := url.QueryEscape(pg.Password)
|
||||
|
||||
// 构建 DSN
|
||||
dsn := fmt.Sprintf(
|
||||
"postgres://%s:%s@%s:%d/%s?sslmode=%s&timezone=%s",
|
||||
pg.User,
|
||||
password,
|
||||
pg.Host,
|
||||
pg.Port,
|
||||
pg.Dbname,
|
||||
sslMode,
|
||||
pg.TimeZone,
|
||||
)
|
||||
|
||||
db, err := pgxpool.New(ctx, dsn)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := db.Ping(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
logger.Log.Info("√ 数据库连接成功")
|
||||
|
||||
return db, nil
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
package bootstrap
|
||||
|
||||
import (
|
||||
"allapp-go-v2/internal/config"
|
||||
"allapp-go-v2/internal/pkg/uniqueid"
|
||||
)
|
||||
|
||||
func InitUniqueId(cfg *config.Config) {
|
||||
options := uniqueid.NewIdGeneratorOptions(cfg.UniqueID.WorkerID)
|
||||
uniqueid.SetIdGenerator(options)
|
||||
}
|
||||
@@ -0,0 +1,103 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
"github.com/spf13/viper"
|
||||
)
|
||||
|
||||
func Load() (*Config, error) {
|
||||
env := getEnv("APP_ENV", "dev")
|
||||
|
||||
v := viper.New()
|
||||
v.SetConfigType("yaml")
|
||||
|
||||
// 配置路径
|
||||
v.AddConfigPath("configs")
|
||||
|
||||
// ======================
|
||||
// 1️⃣ 基础配置
|
||||
// ======================
|
||||
if err := loadBase(v); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// ======================
|
||||
// 2️⃣ 环境配置(支持多环境)
|
||||
// ======================
|
||||
loadByEnv(v, env)
|
||||
|
||||
// ======================
|
||||
// 3️⃣ 环境变量覆盖
|
||||
// ======================
|
||||
bindEnv(v)
|
||||
|
||||
// ======================
|
||||
// 4️⃣ 默认值
|
||||
// ======================
|
||||
setDefaults(v, env)
|
||||
|
||||
// ======================
|
||||
// 5️⃣ 解析
|
||||
// ======================
|
||||
var cfg Config
|
||||
if err := v.Unmarshal(&cfg); err != nil {
|
||||
return nil, fmt.Errorf("decode config failed: %w", err)
|
||||
}
|
||||
|
||||
return &cfg, nil
|
||||
}
|
||||
|
||||
func loadBase(v *viper.Viper) error {
|
||||
v.SetConfigName("config")
|
||||
|
||||
if err := v.ReadInConfig(); err != nil {
|
||||
return fmt.Errorf("read base config failed: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func loadByEnv(v *viper.Viper, env string) {
|
||||
envs := strings.Split(env, ",")
|
||||
|
||||
for _, e := range envs {
|
||||
e = strings.TrimSpace(e)
|
||||
if e == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
v.SetConfigName("config." + e)
|
||||
|
||||
if err := v.MergeInConfig(); err != nil {
|
||||
fmt.Printf("⚠️ config.%s.yaml not found\n", e)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func bindEnv(v *viper.Viper) {
|
||||
v.AutomaticEnv()
|
||||
v.SetEnvKeyReplacer(strings.NewReplacer(".", "_"))
|
||||
}
|
||||
|
||||
func setDefaults(v *viper.Viper, env string) {
|
||||
v.SetDefault("app.env", env)
|
||||
|
||||
v.SetDefault("app.name", "golang")
|
||||
|
||||
v.SetDefault("server.port", "8080")
|
||||
|
||||
v.SetDefault("unique_id.datacenter_id", 1)
|
||||
v.SetDefault("unique_id.worker_id", 1)
|
||||
|
||||
v.SetDefault("jwt.secret", "abcdedhaldkasdlkasd")
|
||||
v.SetDefault("jwt.access_expiry", -1)
|
||||
}
|
||||
|
||||
func getEnv(key, def string) string {
|
||||
if val := os.Getenv(key); val != "" {
|
||||
return val
|
||||
}
|
||||
return def
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
package config
|
||||
|
||||
import "time"
|
||||
|
||||
type Config struct {
|
||||
App AppConfig `mapstructure:"app"`
|
||||
Server ServerConfig `mapstructure:"server"`
|
||||
JWT JwtConfig `mapstructure:"jwt"`
|
||||
UniqueID UniqueIDConfig `mapstructure:"unique_id"`
|
||||
Postgres PostgresConfig `mapstructure:"postgres"`
|
||||
Wechat WechatConfig `mapstructure:"wechat"`
|
||||
Redis RedisConfig `mapstructure:"redis"`
|
||||
}
|
||||
|
||||
type AppConfig struct {
|
||||
Name string `mapstructure:"name"`
|
||||
Env string `mapstructure:"env"`
|
||||
}
|
||||
|
||||
type ServerConfig struct {
|
||||
Port string `mapstructure:"port"`
|
||||
}
|
||||
|
||||
type JwtConfig struct {
|
||||
Secret string `mapstructure:"secret"`
|
||||
AccessExpiry time.Duration `mapstructure:"access_expiry"`
|
||||
}
|
||||
|
||||
type UniqueIDConfig struct {
|
||||
DataCenterID uint16 `mapstructure:"datacenter_id"`
|
||||
WorkerID uint16 `mapstructure:"worker_id"`
|
||||
}
|
||||
|
||||
type WechatConfig struct {
|
||||
AppId string `mapstructure:"app_id"`
|
||||
AppSecret string `mapstructure:"app_secret"`
|
||||
}
|
||||
|
||||
type PostgresConfig struct {
|
||||
Host string `mapstructure:"host"`
|
||||
Port int `mapstructure:"port"`
|
||||
User string `mapstructure:"user"`
|
||||
Password string `mapstructure:"password"`
|
||||
Dbname string `mapstructure:"dbname"`
|
||||
SslMode bool `mapstructure:"ssl_mode"`
|
||||
TimeZone string `mapstructure:"timezone"`
|
||||
}
|
||||
|
||||
type RedisConfig struct {
|
||||
Host string `mapstructure:"host"`
|
||||
Port int `mapstructure:"port"`
|
||||
Db int `mapstructure:"db"`
|
||||
Password string `mapstructure:"password"`
|
||||
}
|
||||
@@ -0,0 +1,147 @@
|
||||
package logger
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
)
|
||||
|
||||
var Log *zap.Logger
|
||||
|
||||
var (
|
||||
currentDate string
|
||||
mu sync.Mutex
|
||||
envMode string
|
||||
)
|
||||
|
||||
// ======================
|
||||
// 初始化
|
||||
// ======================
|
||||
func Init(env string) error {
|
||||
envMode = env
|
||||
logger := buildLogger()
|
||||
|
||||
Log = logger
|
||||
return nil
|
||||
}
|
||||
|
||||
// ======================
|
||||
// 构建 logger
|
||||
// ======================
|
||||
func buildLogger() *zap.Logger {
|
||||
currentDate = getToday()
|
||||
|
||||
level := zap.InfoLevel
|
||||
if envMode == "dev" {
|
||||
level = zap.DebugLevel
|
||||
}
|
||||
|
||||
encoderConfig := zapcore.EncoderConfig{
|
||||
TimeKey: "time",
|
||||
LevelKey: "level",
|
||||
MessageKey: "msg",
|
||||
CallerKey: "caller",
|
||||
EncodeLevel: zapcore.CapitalLevelEncoder,
|
||||
EncodeTime: timeEncoder,
|
||||
EncodeCaller: zapcore.ShortCallerEncoder,
|
||||
}
|
||||
|
||||
encoder := zapcore.NewJSONEncoder(encoderConfig)
|
||||
|
||||
consoleWriter := zapcore.AddSync(os.Stdout)
|
||||
|
||||
infoWriter := zapcore.AddSync(&dailyWriter{
|
||||
level: "info",
|
||||
})
|
||||
|
||||
errorWriter := zapcore.AddSync(&dailyWriter{
|
||||
level: "error",
|
||||
})
|
||||
|
||||
infoCore := zapcore.NewCore(
|
||||
encoder,
|
||||
zapcore.NewMultiWriteSyncer(consoleWriter, infoWriter),
|
||||
zap.LevelEnablerFunc(func(lvl zapcore.Level) bool {
|
||||
return lvl < zapcore.ErrorLevel && lvl >= level
|
||||
}),
|
||||
)
|
||||
|
||||
errorCore := zapcore.NewCore(
|
||||
encoder,
|
||||
zapcore.NewMultiWriteSyncer(consoleWriter, errorWriter),
|
||||
zap.LevelEnablerFunc(func(lvl zapcore.Level) bool {
|
||||
return lvl >= zapcore.ErrorLevel
|
||||
}),
|
||||
)
|
||||
|
||||
core := zapcore.NewTee(infoCore, errorCore)
|
||||
|
||||
return zap.New(core, zap.AddCaller(), zap.AddCallerSkip(1))
|
||||
}
|
||||
|
||||
// ======================
|
||||
// 自定义 Writer(核心)
|
||||
// ======================
|
||||
type dailyWriter struct {
|
||||
level string
|
||||
file *os.File
|
||||
}
|
||||
|
||||
func (w *dailyWriter) Write(p []byte) (n int, err error) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
|
||||
today := getToday()
|
||||
|
||||
// 跨天切换
|
||||
if w.file == nil || today != currentDate {
|
||||
if w.file != nil {
|
||||
_ = w.file.Close()
|
||||
}
|
||||
|
||||
currentDate = today
|
||||
|
||||
filename := getLogFilePath(w.level, today)
|
||||
|
||||
file, err := os.OpenFile(filename, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
w.file = file
|
||||
}
|
||||
|
||||
return w.file.Write(p)
|
||||
}
|
||||
|
||||
func (w *dailyWriter) Sync() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// ======================
|
||||
// 工具函数
|
||||
// ======================
|
||||
|
||||
// 获取当天
|
||||
func getToday() string {
|
||||
return time.Now().Format("2006-01-02")
|
||||
}
|
||||
|
||||
// 获取日志路径
|
||||
func getLogFilePath(level, date string) string {
|
||||
dir := filepath.Join("logs", level)
|
||||
|
||||
// 自动创建目录
|
||||
_ = os.MkdirAll(dir, os.ModePerm)
|
||||
|
||||
return filepath.Join(dir, date+".log")
|
||||
}
|
||||
|
||||
// 时间格式
|
||||
func timeEncoder(t time.Time, enc zapcore.PrimitiveArrayEncoder) {
|
||||
enc.AppendString(t.Format("2006-01-02 15:04:05"))
|
||||
}
|
||||
@@ -0,0 +1,83 @@
|
||||
package uniqueid
|
||||
|
||||
import (
|
||||
"strconv"
|
||||
"time"
|
||||
)
|
||||
|
||||
type DefaultIdGenerator struct {
|
||||
Options *IdGeneratorOptions
|
||||
SnowWorker ISnowWorker
|
||||
IdGeneratorException IdGeneratorException
|
||||
}
|
||||
|
||||
func NewDefaultIdGenerator(options *IdGeneratorOptions) *DefaultIdGenerator {
|
||||
if options == nil {
|
||||
panic("dig.Options error.")
|
||||
}
|
||||
|
||||
// 1.BaseTime
|
||||
minTime := int64(631123200000) // time.Now().AddDate(-30, 0, 0).UnixNano() / 1e6
|
||||
if options.BaseTime < minTime || options.BaseTime > time.Now().UnixNano()/1e6 {
|
||||
panic("BaseTime error.")
|
||||
}
|
||||
|
||||
// 2.WorkerIdBitLength
|
||||
if options.WorkerIdBitLength <= 0 {
|
||||
panic("WorkerIdBitLength error.(range:[1, 21])")
|
||||
}
|
||||
if options.WorkerIdBitLength+options.SeqBitLength > 22 {
|
||||
panic("error:WorkerIdBitLength + SeqBitLength <= 22")
|
||||
}
|
||||
|
||||
// 3.WorkerId
|
||||
maxWorkerIdNumber := uint16(1<<options.WorkerIdBitLength) - 1
|
||||
if maxWorkerIdNumber == 0 {
|
||||
maxWorkerIdNumber = 63
|
||||
}
|
||||
if options.WorkerId < 0 || options.WorkerId > maxWorkerIdNumber {
|
||||
panic("WorkerId error. (range:[0, " + strconv.FormatUint(uint64(maxWorkerIdNumber), 10) + "]")
|
||||
}
|
||||
|
||||
// 4.SeqBitLength
|
||||
if options.SeqBitLength < 2 || options.SeqBitLength > 21 {
|
||||
panic("SeqBitLength error. (range:[2, 21])")
|
||||
}
|
||||
|
||||
// 5.MaxSeqNumber
|
||||
maxSeqNumber := uint32(1<<options.SeqBitLength) - 1
|
||||
if maxSeqNumber == 0 {
|
||||
maxSeqNumber = 63
|
||||
}
|
||||
if options.MaxSeqNumber < 0 || options.MaxSeqNumber > maxSeqNumber {
|
||||
panic("MaxSeqNumber error. (range:[1, " + strconv.FormatUint(uint64(maxSeqNumber), 10) + "]")
|
||||
}
|
||||
|
||||
// 6.MinSeqNumber
|
||||
if options.MinSeqNumber < 5 || options.MinSeqNumber > maxSeqNumber {
|
||||
panic("MinSeqNumber error. (range:[5, " + strconv.FormatUint(uint64(maxSeqNumber), 10) + "]")
|
||||
}
|
||||
|
||||
var snowWorker ISnowWorker
|
||||
switch options.Method {
|
||||
case 1:
|
||||
snowWorker = NewSnowWorkerM1(options)
|
||||
case 2:
|
||||
snowWorker = NewSnowWorkerM2(options)
|
||||
default:
|
||||
snowWorker = NewSnowWorkerM1(options)
|
||||
}
|
||||
|
||||
if options.Method == 1 {
|
||||
time.Sleep(time.Duration(500) * time.Microsecond)
|
||||
}
|
||||
|
||||
return &DefaultIdGenerator{
|
||||
Options: options,
|
||||
SnowWorker: snowWorker,
|
||||
}
|
||||
}
|
||||
|
||||
func (dig DefaultIdGenerator) NewLong() int64 {
|
||||
return dig.SnowWorker.NextId()
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
package uniqueid
|
||||
|
||||
type IIdGenerator interface {
|
||||
NewLong() uint64
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
package uniqueid
|
||||
|
||||
type ISnowWorker interface {
|
||||
NextId() int64
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package uniqueid
|
||||
|
||||
import "fmt"
|
||||
|
||||
type IdGeneratorException struct {
|
||||
message string
|
||||
error error
|
||||
}
|
||||
|
||||
func (e IdGeneratorException) IdGeneratorException(message ...interface{}) {
|
||||
fmt.Println(message)
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package uniqueid
|
||||
|
||||
type IdGeneratorOptions struct {
|
||||
Method uint16 // 雪花计算方法,(1-漂移算法|2-传统算法),默认1
|
||||
BaseTime int64 // 基础时间(ms单位),不能超过当前系统时间
|
||||
WorkerId uint16 // 机器码,必须由外部设定,最大值 2^WorkerIdBitLength-1
|
||||
WorkerIdBitLength byte // 机器码位长,默认值6,取值范围 [1, 15](要求:序列数位长+机器码位长不超过22)
|
||||
SeqBitLength byte // 序列数位长,默认值6,取值范围 [3, 21](要求:序列数位长+机器码位长不超过22)
|
||||
MaxSeqNumber uint32 // 最大序列数(含),设置范围 [MinSeqNumber, 2^SeqBitLength-1],默认值0,表示最大序列数取最大值(2^SeqBitLength-1])
|
||||
MinSeqNumber uint32 // 最小序列数(含),默认值5,取值范围 [5, MaxSeqNumber],每毫秒的前5个序列数对应编号0-4是保留位,其中1-4是时间回拨相应预留位,0是手工新值预留位
|
||||
TopOverCostCount uint32 // 最大漂移次数(含),默认2000,推荐范围500-10000(与计算能力有关)
|
||||
}
|
||||
|
||||
func NewIdGeneratorOptions(workerId uint16) *IdGeneratorOptions {
|
||||
return &IdGeneratorOptions{
|
||||
Method: 1,
|
||||
WorkerId: workerId,
|
||||
BaseTime: 1582136402000,
|
||||
WorkerIdBitLength: 6,
|
||||
SeqBitLength: 6,
|
||||
MaxSeqNumber: 0,
|
||||
MinSeqNumber: 5,
|
||||
TopOverCostCount: 2000,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
package uniqueid
|
||||
|
||||
import (
|
||||
"sync"
|
||||
)
|
||||
|
||||
var singletonMutex sync.Mutex
|
||||
var idGenerator *DefaultIdGenerator
|
||||
|
||||
// SetIdGenerator .
|
||||
func SetIdGenerator(options *IdGeneratorOptions) {
|
||||
singletonMutex.Lock()
|
||||
idGenerator = NewDefaultIdGenerator(options)
|
||||
singletonMutex.Unlock()
|
||||
}
|
||||
|
||||
// NextId .
|
||||
func NextId() int64 {
|
||||
if idGenerator == nil {
|
||||
singletonMutex.Lock()
|
||||
defer singletonMutex.Unlock()
|
||||
if idGenerator == nil {
|
||||
options := NewIdGeneratorOptions(1)
|
||||
idGenerator = NewDefaultIdGenerator(options)
|
||||
}
|
||||
}
|
||||
|
||||
return idGenerator.NewLong()
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
package uniqueid
|
||||
|
||||
type OverCostActionArg struct {
|
||||
ActionType int32
|
||||
TimeTick int64
|
||||
WorkerId uint16
|
||||
OverCostCountInOneTerm int32
|
||||
GenCountInOneTerm int32
|
||||
TermIndex int32
|
||||
}
|
||||
|
||||
func (ocaa OverCostActionArg) OverCostActionArg(workerId uint16, timeTick int64, actionType int32, overCostCountInOneTerm int32, genCountWhenOverCost int32, index int32) {
|
||||
ocaa.ActionType = actionType
|
||||
ocaa.TimeTick = timeTick
|
||||
ocaa.WorkerId = workerId
|
||||
ocaa.OverCostCountInOneTerm = overCostCountInOneTerm
|
||||
ocaa.GenCountInOneTerm = genCountWhenOverCost
|
||||
ocaa.TermIndex = index
|
||||
}
|
||||
@@ -0,0 +1,243 @@
|
||||
package uniqueid
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// SnowWorkerM1 .
|
||||
type SnowWorkerM1 struct {
|
||||
BaseTime int64 //基础时间
|
||||
WorkerId uint16 //机器码
|
||||
WorkerIdBitLength byte //机器码位长
|
||||
SeqBitLength byte //自增序列数位长
|
||||
MaxSeqNumber uint32 //最大序列数(含)
|
||||
MinSeqNumber uint32 //最小序列数(含)
|
||||
TopOverCostCount uint32 //最大漂移次数
|
||||
_TimestampShift byte
|
||||
_CurrentSeqNumber uint32
|
||||
|
||||
_LastTimeTick int64
|
||||
_TurnBackTimeTick int64
|
||||
_TurnBackIndex byte
|
||||
_IsOverCost bool
|
||||
_OverCostCountInOneTerm uint32
|
||||
_GenCountInOneTerm uint32
|
||||
_TermIndex uint32
|
||||
|
||||
sync.Mutex
|
||||
}
|
||||
|
||||
// NewSnowWorkerM1 .
|
||||
func NewSnowWorkerM1(options *IdGeneratorOptions) ISnowWorker {
|
||||
var workerIdBitLength byte
|
||||
var seqBitLength byte
|
||||
var maxSeqNumber uint32
|
||||
|
||||
// 1.BaseTime
|
||||
var baseTime int64
|
||||
if options.BaseTime != 0 {
|
||||
baseTime = options.BaseTime
|
||||
} else {
|
||||
baseTime = 1582136402000
|
||||
}
|
||||
|
||||
// 2.WorkerIdBitLength
|
||||
if options.WorkerIdBitLength == 0 {
|
||||
workerIdBitLength = 6
|
||||
} else {
|
||||
workerIdBitLength = options.WorkerIdBitLength
|
||||
}
|
||||
|
||||
// 3.WorkerId
|
||||
var workerId = options.WorkerId
|
||||
|
||||
// 4.SeqBitLength
|
||||
if options.SeqBitLength == 0 {
|
||||
seqBitLength = 6
|
||||
} else {
|
||||
seqBitLength = options.SeqBitLength
|
||||
}
|
||||
|
||||
// 5.MaxSeqNumber
|
||||
if options.MaxSeqNumber <= 0 {
|
||||
maxSeqNumber = (1 << seqBitLength) - 1
|
||||
} else {
|
||||
maxSeqNumber = options.MaxSeqNumber
|
||||
}
|
||||
|
||||
// 6.MinSeqNumber
|
||||
var minSeqNumber = options.MinSeqNumber
|
||||
|
||||
// 7.Others
|
||||
var topOverCostCount = options.TopOverCostCount
|
||||
if topOverCostCount == 0 {
|
||||
topOverCostCount = 2000
|
||||
}
|
||||
|
||||
timestampShift := (byte)(workerIdBitLength + seqBitLength)
|
||||
currentSeqNumber := minSeqNumber
|
||||
|
||||
return &SnowWorkerM1{
|
||||
BaseTime: baseTime,
|
||||
WorkerIdBitLength: workerIdBitLength,
|
||||
WorkerId: workerId,
|
||||
SeqBitLength: seqBitLength,
|
||||
MaxSeqNumber: maxSeqNumber,
|
||||
MinSeqNumber: minSeqNumber,
|
||||
TopOverCostCount: topOverCostCount,
|
||||
_TimestampShift: timestampShift,
|
||||
_CurrentSeqNumber: currentSeqNumber,
|
||||
|
||||
_LastTimeTick: 0,
|
||||
_TurnBackTimeTick: 0,
|
||||
_TurnBackIndex: 0,
|
||||
_IsOverCost: false,
|
||||
_OverCostCountInOneTerm: 0,
|
||||
_GenCountInOneTerm: 0,
|
||||
_TermIndex: 0,
|
||||
}
|
||||
}
|
||||
|
||||
// DoGenIDAction .
|
||||
func (m1 *SnowWorkerM1) DoGenIdAction(arg *OverCostActionArg) {
|
||||
|
||||
}
|
||||
|
||||
func (m1 *SnowWorkerM1) BeginOverCostAction(useTimeTick int64) {
|
||||
|
||||
}
|
||||
|
||||
func (m1 *SnowWorkerM1) EndOverCostAction(useTimeTick int64) {
|
||||
if m1._TermIndex > 10000 {
|
||||
m1._TermIndex = 0
|
||||
}
|
||||
}
|
||||
|
||||
func (m1 *SnowWorkerM1) BeginTurnBackAction(useTimeTick int64) {
|
||||
|
||||
}
|
||||
|
||||
func (m1 *SnowWorkerM1) EndTurnBackAction(useTimeTick int64) {
|
||||
|
||||
}
|
||||
|
||||
func (m1 *SnowWorkerM1) NextOverCostId() int64 {
|
||||
currentTimeTick := m1.GetCurrentTimeTick()
|
||||
if currentTimeTick > m1._LastTimeTick {
|
||||
m1.EndOverCostAction(currentTimeTick)
|
||||
m1._LastTimeTick = currentTimeTick
|
||||
m1._CurrentSeqNumber = m1.MinSeqNumber
|
||||
m1._IsOverCost = false
|
||||
m1._OverCostCountInOneTerm = 0
|
||||
m1._GenCountInOneTerm = 0
|
||||
return m1.CalcId(m1._LastTimeTick)
|
||||
}
|
||||
if m1._OverCostCountInOneTerm >= m1.TopOverCostCount {
|
||||
m1.EndOverCostAction(currentTimeTick)
|
||||
m1._LastTimeTick = m1.GetNextTimeTick()
|
||||
m1._CurrentSeqNumber = m1.MinSeqNumber
|
||||
m1._IsOverCost = false
|
||||
m1._OverCostCountInOneTerm = 0
|
||||
m1._GenCountInOneTerm = 0
|
||||
return m1.CalcId(m1._LastTimeTick)
|
||||
}
|
||||
if m1._CurrentSeqNumber > m1.MaxSeqNumber {
|
||||
m1._LastTimeTick++
|
||||
m1._CurrentSeqNumber = m1.MinSeqNumber
|
||||
m1._IsOverCost = true
|
||||
m1._OverCostCountInOneTerm++
|
||||
m1._GenCountInOneTerm++
|
||||
|
||||
return m1.CalcId(m1._LastTimeTick)
|
||||
}
|
||||
|
||||
m1._GenCountInOneTerm++
|
||||
return m1.CalcId(m1._LastTimeTick)
|
||||
}
|
||||
|
||||
// NextNormalID .
|
||||
func (m1 *SnowWorkerM1) NextNormalId() int64 {
|
||||
currentTimeTick := m1.GetCurrentTimeTick()
|
||||
if currentTimeTick < m1._LastTimeTick {
|
||||
if m1._TurnBackTimeTick < 1 {
|
||||
m1._TurnBackTimeTick = m1._LastTimeTick - 1
|
||||
m1._TurnBackIndex++
|
||||
// 每毫秒序列数的前5位是预留位,0用于手工新值,1-4是时间回拨次序
|
||||
// 最多4次回拨(防止回拨重叠)
|
||||
if m1._TurnBackIndex > 4 {
|
||||
m1._TurnBackIndex = 1
|
||||
}
|
||||
m1.BeginTurnBackAction(m1._TurnBackTimeTick)
|
||||
}
|
||||
|
||||
// time.Sleep(time.Duration(1) * time.Millisecond)
|
||||
return m1.CalcTurnBackId(m1._TurnBackTimeTick)
|
||||
}
|
||||
|
||||
// 时间追平时,_TurnBackTimeTick清零
|
||||
if m1._TurnBackTimeTick > 0 {
|
||||
m1.EndTurnBackAction(m1._TurnBackTimeTick)
|
||||
m1._TurnBackTimeTick = 0
|
||||
}
|
||||
|
||||
if currentTimeTick > m1._LastTimeTick {
|
||||
m1._LastTimeTick = currentTimeTick
|
||||
m1._CurrentSeqNumber = m1.MinSeqNumber
|
||||
return m1.CalcId(m1._LastTimeTick)
|
||||
}
|
||||
|
||||
if m1._CurrentSeqNumber > m1.MaxSeqNumber {
|
||||
m1.BeginOverCostAction(currentTimeTick)
|
||||
m1._TermIndex++
|
||||
m1._LastTimeTick++
|
||||
m1._CurrentSeqNumber = m1.MinSeqNumber
|
||||
m1._IsOverCost = true
|
||||
m1._OverCostCountInOneTerm = 1
|
||||
m1._GenCountInOneTerm = 1
|
||||
|
||||
return m1.CalcId(m1._LastTimeTick)
|
||||
}
|
||||
|
||||
return m1.CalcId(m1._LastTimeTick)
|
||||
}
|
||||
|
||||
// CalcID .
|
||||
func (m1 *SnowWorkerM1) CalcId(useTimeTick int64) int64 {
|
||||
result := int64(useTimeTick<<m1._TimestampShift) + int64(m1.WorkerId<<m1.SeqBitLength) + int64(m1._CurrentSeqNumber)
|
||||
m1._CurrentSeqNumber++
|
||||
return result
|
||||
}
|
||||
|
||||
// CalcTurnBackID .
|
||||
func (m1 *SnowWorkerM1) CalcTurnBackId(useTimeTick int64) int64 {
|
||||
result := int64(useTimeTick<<m1._TimestampShift) + int64(m1.WorkerId<<m1.SeqBitLength) + int64(m1._TurnBackIndex)
|
||||
m1._TurnBackTimeTick--
|
||||
return result
|
||||
}
|
||||
|
||||
// GetCurrentTimeTick .
|
||||
func (m1 *SnowWorkerM1) GetCurrentTimeTick() int64 {
|
||||
var millis = time.Now().UnixNano() / 1e6
|
||||
return millis - m1.BaseTime
|
||||
}
|
||||
|
||||
// GetNextTimeTick .
|
||||
func (m1 *SnowWorkerM1) GetNextTimeTick() int64 {
|
||||
tempTimeTicker := m1.GetCurrentTimeTick()
|
||||
for tempTimeTicker <= m1._LastTimeTick {
|
||||
tempTimeTicker = m1.GetCurrentTimeTick()
|
||||
}
|
||||
return tempTimeTicker
|
||||
}
|
||||
|
||||
// NextId .
|
||||
func (m1 *SnowWorkerM1) NextId() int64 {
|
||||
m1.Lock()
|
||||
defer m1.Unlock()
|
||||
if m1._IsOverCost {
|
||||
return m1.NextOverCostId()
|
||||
} else {
|
||||
return m1.NextNormalId()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
package uniqueid
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
)
|
||||
|
||||
type SnowWorkerM2 struct {
|
||||
*SnowWorkerM1
|
||||
}
|
||||
|
||||
func NewSnowWorkerM2(options *IdGeneratorOptions) ISnowWorker {
|
||||
return &SnowWorkerM2{
|
||||
NewSnowWorkerM1(options).(*SnowWorkerM1),
|
||||
}
|
||||
}
|
||||
|
||||
func (m2 SnowWorkerM2) NextId() int64 {
|
||||
m2.Lock()
|
||||
defer m2.Unlock()
|
||||
currentTimeTick := m2.GetCurrentTimeTick()
|
||||
if m2._LastTimeTick == currentTimeTick {
|
||||
m2._CurrentSeqNumber++
|
||||
if m2._CurrentSeqNumber > m2.MaxSeqNumber {
|
||||
m2._CurrentSeqNumber = m2.MinSeqNumber
|
||||
currentTimeTick = m2.GetNextTimeTick()
|
||||
}
|
||||
} else {
|
||||
m2._CurrentSeqNumber = m2.MinSeqNumber
|
||||
}
|
||||
if currentTimeTick < m2._LastTimeTick {
|
||||
fmt.Println("Time error for {0} milliseconds", strconv.FormatInt(m2._LastTimeTick-currentTimeTick, 10))
|
||||
}
|
||||
m2._LastTimeTick = currentTimeTick
|
||||
result := int64(currentTimeTick<<m2._TimestampShift) + int64(m2.WorkerId<<m2.SeqBitLength) + int64(m2._CurrentSeqNumber)
|
||||
return result
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"allapp-go-v2/internal/bootstrap"
|
||||
"fmt"
|
||||
|
||||
"github.com/gofiber/fiber/v3"
|
||||
"github.com/gofiber/fiber/v3/middleware/requestid"
|
||||
)
|
||||
|
||||
func Start(appCtx *bootstrap.App) error {
|
||||
// 1️⃣ 创建 fiber 实例
|
||||
app := fiber.New()
|
||||
|
||||
app.Use(requestid.New())
|
||||
|
||||
// 2️⃣ 注册路由
|
||||
//registerRoutes(app, appCtx)
|
||||
|
||||
// 3️⃣ 启动服务
|
||||
addr := fmt.Sprintf(":%s", appCtx.Config.Server.Port)
|
||||
|
||||
fmt.Println("🚀 server running at", addr)
|
||||
|
||||
return app.Listen(addr)
|
||||
}
|
||||
Reference in new issue
Block a user