pipeline-convert/main.go

105 lines
2.5 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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")
}