forked from ci4s/pipeline-convert
105 lines
2.5 KiB
Go
105 lines
2.5 KiB
Go
package main
|
||
|
||
import (
|
||
"github.com/go-redsync/redsync/v4"
|
||
"os"
|
||
"os/signal"
|
||
"pipeline-convert/configs"
|
||
"pipeline-convert/internal/dao"
|
||
"pipeline-convert/internal/handler"
|
||
"pipeline-convert/internal/log"
|
||
"pipeline-convert/internal/router"
|
||
"pipeline-convert/pkg/k8slistener"
|
||
"pipeline-convert/pkg/redislock"
|
||
"syscall"
|
||
"time"
|
||
)
|
||
|
||
// @title 流水线转换服务接口文档
|
||
// @description 流水线转换服务接口文档.
|
||
// @version 1.0
|
||
// @host 开发环境:172.20.32.181:31000,测试环境:172.20.32.185:31000
|
||
// @BasePath /
|
||
func main() {
|
||
configs.Init()
|
||
dao.Init()
|
||
r := router.RegisterRouter()
|
||
r.Use(log.GinLogger())
|
||
startListen()
|
||
r.Run(":80")
|
||
}
|
||
|
||
func startListen() {
|
||
redisConfig := configs.RConfig.GetConfig().RedisConfig
|
||
lock := redislock.NewRedlockClient(&redisConfig)
|
||
failed := 0
|
||
var mutex *redsync.Mutex
|
||
var err error
|
||
for {
|
||
mutex, err = lock.Lock("pipeline-convert", time.Second*10)
|
||
if err != nil {
|
||
failed++
|
||
if failed > 3 {
|
||
log.Error("redis lock acquire failed too many times, exit")
|
||
return
|
||
}
|
||
} else {
|
||
break
|
||
}
|
||
time.Sleep(5 * time.Second)
|
||
}
|
||
stopCh := make(chan os.Signal, 1)
|
||
signal.Notify(stopCh, syscall.SIGINT, syscall.SIGTERM)
|
||
log.Info("redis lock acquire success")
|
||
|
||
go func() {
|
||
// 续期锁
|
||
failed = 0
|
||
// 使用定时器续期锁
|
||
ticker := time.NewTicker(5 * time.Second)
|
||
defer ticker.Stop() // 确保在程序结束时停止Ticker
|
||
// 续期失败次数超过3次,退出程序
|
||
for {
|
||
select {
|
||
case <-ticker.C:
|
||
ttl := mutex.Until()
|
||
log.Debugf("redis lock extend start, expirtion time: %v", ttl)
|
||
|
||
ret, err := mutex.Extend()
|
||
if !ret || (err != nil) {
|
||
log.Errorf("redis lock extend failed, err: %v", err)
|
||
failed++
|
||
} else {
|
||
failed = 0
|
||
log.Debug("redis lock extend success")
|
||
}
|
||
if failed > 3 {
|
||
log.Error("redis lock extend failed too many times, exit")
|
||
return
|
||
}
|
||
case <-stopCh:
|
||
log.Info("receive stop signal, release redis lock")
|
||
if ret, err := mutex.Unlock(); !ret || (err != nil) {
|
||
log.Error("redis lock release failed, ret: %v, err: %v", ret, err)
|
||
}
|
||
return
|
||
}
|
||
}
|
||
}()
|
||
|
||
client, err := handler.GetK8sClient()
|
||
if err != nil {
|
||
log.Errorf("create k8s client failed, err: %v", err)
|
||
return
|
||
}
|
||
// start listen
|
||
listen, err := k8slistener.NewListener(client, "app-deploy=model-service", &redisConfig)
|
||
if err != nil {
|
||
log.Errorf("create k8s listener failed, err: %v", err)
|
||
return
|
||
}
|
||
ch := make(chan struct{})
|
||
go listen.Start(ch)
|
||
log.Info("start listen success")
|
||
}
|