diff --git a/config.go b/config.go index d567cbe..276ff87 100644 --- a/config.go +++ b/config.go @@ -3,8 +3,6 @@ package main import ( "encoding/json" "os" - "sync" - "time" ) // Config 对应 config.json 结构 @@ -16,11 +14,6 @@ type Config struct { var ( globalConfig Config - // lastExecMap 冷却去重状态,key: "repoName:ref", value: *debounceEntry - // 记录每个 key 上次实际执行时刻,冷却窗口内到达的 push 直接丢弃 - lastExecMap sync.Map - // 防抖冷却时间限制:3 分钟 - debounceDuration = 3 * time.Minute // bark 推送消息 pushURL = "https://bark.maimaicuizhiji.top/push" ) diff --git a/main.go b/main.go index a41b46a..30c961d 100644 --- a/main.go +++ b/main.go @@ -27,10 +27,12 @@ func main() { log.Printf("[%s] %s %d %v", c.Request.Method, c.Request.URL.Path, c.Writer.Status(), time.Since(t)) }) + startBuildWorker() + r.POST("/webhook", handleWebhook) listenAddr := fmt.Sprintf(":%d", *portFlag) - log.Printf("GitLab Webhook 服务已启动 (3 分钟防抖去重),监听端口 %d ...\n", *portFlag) + log.Printf("GitLab Webhook 服务已启动 (队列模式 + 3 分钟冷却去重),监听端口 %d ...\n", *portFlag) if err := r.Run(listenAddr); err != nil { log.Fatalf("服务启动失败: %v", err) diff --git a/runner.go b/runner.go index 407ea9a..2e99193 100644 --- a/runner.go +++ b/runner.go @@ -17,6 +17,77 @@ import ( // scriptTimeout 脚本执行超时上限(15 分钟) const scriptTimeout = 60 * time.Minute +// buildTask 单次构建任务的快照,handler 入队时构造。 +type buildTask struct { + scriptPath string + repoName string + ref string + commitID string +} + +// debounceEntry 单 key 的冷却去重状态。 +type debounceEntry struct { + mu sync.Mutex + lastExec time.Time +} + +var ( + // buildCh 单槽队列;handler enqueue 始终非阻塞,worker 唯一消费者 + buildCh = make(chan buildTask, 1) + // lastExecMap 冷却去重状态,key: "repoName:ref", value: *debounceEntry + lastExecMap sync.Map + // 防抖冷却时间限制:3 分钟 + debounceDuration = 3 * time.Minute +) + +// enqueueBuildTask 入队;队列已满时丢老取新(drop-oldest) +func enqueueBuildTask(t buildTask) { + select { + case buildCh <- t: + return // 入队成功 + default: + // 队列已满,丢弃旧 task + select { + case <-buildCh: + default: + } + buildCh <- t // 阻塞不会发生,因为我们是唯一发送方 + } +} + +// startBuildWorker 启动单 worker goroutine 持续消费 buildCh +func startBuildWorker() { + go consumeBuild() +} + +// consumeBuild worker 主循环;取出 task,按冷却规则执行或跳过 +func consumeBuild() { + for task := range buildCh { + lockKey := fmt.Sprintf("%s:%s", task.repoName, task.ref) + entryVal, _ := lastExecMap.LoadOrStore(lockKey, &debounceEntry{}) + entry := entryVal.(*debounceEntry) + + entry.mu.Lock() + now := time.Now() + if !entry.lastExec.IsZero() && now.Sub(entry.lastExec) < debounceDuration { + remainingSeconds := int((debounceDuration - now.Sub(entry.lastExec)).Seconds()) + if remainingSeconds < 0 { + remainingSeconds = 0 + } + entry.mu.Unlock() + log.Printf("[SKIP] 3 分钟冷却中,跳过构建 [%s] commit %s 剩余 %d 秒\n", + lockKey, shortCommit(task.commitID), remainingSeconds) + continue + } + entry.lastExec = now + entry.mu.Unlock() + + log.Printf("[INFO] 开始执行脚本: %s (仓库: %s, 分支: %s, commit: %s)\n", + task.scriptPath, task.repoName, task.ref, shortCommit(task.commitID)) + go runScript(task.scriptPath, task.repoName, task.ref, task.commitID) + } +} + // shortCommit 安全地取 commit 短 ID,避免空值/短值切片越界 panic func shortCommit(commitID string) string { if len(commitID) > 7 { diff --git a/webhook.go b/webhook.go index d98cdec..5119ee6 100644 --- a/webhook.go +++ b/webhook.go @@ -5,8 +5,6 @@ import ( "log" "net/http" "os" - "sync" - "time" "github.com/gin-gonic/gin" ) @@ -20,13 +18,6 @@ type GitLabPayload struct { } `json:"project"` } -// debounceEntry 单个 key 的冷却去重状态。 -// lastExec:上次实际执行时刻;冷却窗口内到达的 push 直接丢弃,无补跑。 -type debounceEntry struct { - mu sync.Mutex - lastExec time.Time -} - func handleWebhook(c *gin.Context) { // 1. 校验 GitLab Secret Token Header clientToken := c.GetHeader("X-Gitlab-Token") @@ -77,37 +68,18 @@ func handleWebhook(c *gin.Context) { return } - // 6. 冷却去重:窗口外 push 立即执行并刷新 lastExec,窗口内 push 直接丢弃 - lockKey := fmt.Sprintf("%s:%s", repoName, ref) - entryVal, _ := lastExecMap.LoadOrStore(lockKey, &debounceEntry{}) - entry := entryVal.(*debounceEntry) + // 6. 入队等待 worker 按队列顺序消费 + enqueueBuildTask(buildTask{ + scriptPath: scriptPath, + repoName: repoName, + ref: ref, + commitID: commitID, + }) - entry.mu.Lock() - now := time.Now() - - if !entry.lastExec.IsZero() && now.Sub(entry.lastExec) < debounceDuration { - remainingSeconds := max(int((debounceDuration - now.Sub(entry.lastExec)).Seconds()), 0) - entry.mu.Unlock() - - msg := fmt.Sprintf("3 分钟内频繁提交被拦截, 剩余冷却时间: %d 秒", remainingSeconds) - log.Printf("[DEBOUNCE] 跳过执行 [%s] %s\n", lockKey, msg) - c.JSON(http.StatusOK, gin.H{ - "status": "debounced", - "message": msg, - "remaining_seconds": remainingSeconds, - }) - return - } - - entry.lastExec = now - entry.mu.Unlock() - - log.Printf("[INFO] 开始执行脚本: %s (仓库: %s, 分支: %s, commit: %s)\n", + log.Printf("[QUEUED] 收到 push 入队 [%s] (仓库: %s, 分支: %s, commit: %s)\n", scriptPath, repoName, ref, shortCommit(commitID)) - go runScript(scriptPath, repoName, ref, commitID) - c.JSON(http.StatusOK, gin.H{ - "status": "triggered", + "status": "queued", "repository": repoName, "ref": ref, "script": scriptPath,