Files
superwork/tui/internal/tui/webhook.go
T
2026-06-23 05:02:15 +08:00

487 lines
14 KiB
Go

package tui
import (
"context"
"fmt"
"net"
"os"
"path/filepath"
"regexp"
"strconv"
"strings"
"time"
tea "charm.land/bubbletea/v2"
"superwork-tui/internal/auth"
"superwork-tui/internal/cc"
"superwork-tui/internal/config"
"superwork-tui/internal/git"
"superwork-tui/internal/gitea"
"superwork-tui/internal/issue"
"superwork-tui/internal/logging"
"superwork-tui/internal/webhook"
)
// webhookEventMsg carries a single parsed webhook event into the Bubble Tea loop.
type webhookEventMsg struct{ ev webhook.Event }
// webhookServerErrMsg is produced when the webhook server fails to start.
type webhookServerErrMsg struct{ err error }
// autoReviewCapturedMsg carries the codex reviewSessionId captured by the watcher.
type autoReviewCapturedMsg struct {
issueNumber int
sessionID string
}
// autoReviewPersistMsg is produced after writing reviewSessionId to state JSON.
type autoReviewPersistMsg struct {
issueNumber int
err error
}
// configLoadFn is injectable for tests.
var configLoadFn = config.Load
// branchSyncCmdFn is injectable for tests.
var branchSyncCmdFn = branchSyncCheckCmd
// webhookWatchCodexSessionFn is injectable for tests.
var webhookWatchCodexSessionFn = func(ctx context.Context, opts cc.CodexSessionWatchOpts) (string, error) {
return cc.WatchForNewCodexSession(ctx, opts)
}
var reviewMarkerStripRe = regexp.MustCompile(`(?i)<!--\s*spx:review=1\s*-->\s*\n?`)
// specMarkerRe and planMarkerRe match coordinator.ts handleIssueEdited marker extraction.
// Path must contain "/" and end with ".md" to avoid placeholder values.
var specMarkerRe = regexp.MustCompile(`<!--\s*spx:spec=([^\s>]*\/[^\s>]*\.md)\s*-->`)
var planMarkerRe = regexp.MustCompile(`<!--\s*spx:plan=([^\s>]*\/[^\s>]*\.md)\s*-->`)
// waitForWebhookCmd blocks until one event arrives on ch, then delivers it.
// The caller must re-issue this cmd after handling each event to keep listening.
func waitForWebhookCmd(ctx context.Context, ch chan webhook.Event) tea.Cmd {
return func() tea.Msg {
select {
case ev := <-ch:
return webhookEventMsg{ev: ev}
case <-ctx.Done():
return nil
}
}
}
// startWebhookServer starts the HTTP server in a goroutine and returns a Cmd
// that delivers webhookServerErrMsg on port-in-use, or nil msg on success.
// The server stops when ctx is cancelled.
func startWebhookServer(ch chan webhook.Event, ctx context.Context) tea.Cmd {
return func() tea.Msg {
settings, err := configLoadFn()
port := config.DefaultWebhookPort
if err == nil && settings.WebhookPort > 0 {
port = settings.WebhookPort
}
addr := fmt.Sprintf(":%d", port)
srv := webhook.NewServer(func(ev webhook.Event) {
select {
case ch <- ev:
default: // drop if buffer full; non-fatal
}
})
// Probe port availability before blocking.
l, probeErr := net.Listen("tcp", addr)
if probeErr != nil {
return webhookServerErrMsg{err: fmt.Errorf("端口 %d 已被占用", port)}
}
l.Close()
go func() {
if startErr := srv.Start(ctx, addr); startErr != nil && ctx.Err() == nil {
// Ignore ErrServerClosed — normal on shutdown.
_ = startErr
}
}()
return nil
}
}
// dispatchWebhookEvent handles a single webhook event; always re-issues waitForWebhookCmd.
func dispatchWebhookEvent(m Model, ev webhook.Event) (Model, tea.Cmd) {
var cmds []tea.Cmd
cmds = append(cmds, waitForWebhookCmd(m.webhookCtx, m.webhookCh))
logging.Default.Info("webhook", fmt.Sprintf("收到事件 %T", ev))
switch e := ev.(type) {
case webhook.PrEvent:
m, cmds = handlePrEvent(m, e, cmds)
case webhook.IssueEvent:
m, cmds = handleIssueEvent(m, e, cmds)
case webhook.IssueCommentEvent:
m, cmds = handleIssueCommentEvent(m, e, cmds)
case webhook.PushEvent:
m, cmds = handlePushEvent(m, e, cmds)
}
return m, tea.Batch(cmds...)
}
func handlePrEvent(m Model, ev webhook.PrEvent, cmds []tea.Cmd) (Model, []tea.Cmd) {
switch ev.Action {
case "opened", "reopened":
issNum, ok := webhook.ResolveIssueNumber(ev.Body, ev.Branch, ev.IssueNumber)
if !ok {
m.statusMsg = fmt.Sprintf("webhook PR opened: cannot resolve issue# (branch=%s)", ev.Branch)
return m, cmds
}
cmds = append(cmds,
mergePRStateCmd(issNum, ev.PR, ev.Branch),
conditionalAutoReviewCmd(issNum, ev.PR, m.issues),
loadCmd(),
)
m.statusMsg = fmt.Sprintf("PR !%s opened for #%d", ev.PR, issNum)
case "synchronize", "synchronized":
for _, iss := range m.issues {
if iss.Branch == ev.Branch && iss.ReviewSessionID != "" {
cmds = append(cmds, triggerAutoReviewCmd(iss.Number, ev.PR, iss.WorktreePath))
m.statusMsg = fmt.Sprintf("PR !%s synchronize → 更新审查 #%d", ev.PR, iss.Number)
break
}
}
cmds = append(cmds, loadCmd())
case "closed":
issNum, ok := webhook.ResolveIssueNumber(ev.Body, ev.Branch, ev.IssueNumber)
if !ok {
m.statusMsg = fmt.Sprintf("webhook PR closed: cannot resolve issue# (branch=%s)", ev.Branch)
return m, cmds
}
cmds = append(cmds, prClosedCmd(issNum, ev.PR), loadCmd())
m.statusMsg = fmt.Sprintf("PR !%s closed for #%d", ev.PR, issNum)
case "deleted":
issNum, ok := webhook.ResolveIssueNumber(ev.Body, ev.Branch, ev.IssueNumber)
if !ok {
m.statusMsg = fmt.Sprintf("webhook PR deleted: cannot resolve issue# (branch=%s)", ev.Branch)
return m, cmds
}
cmds = append(cmds, prDeletedCmd(issNum), loadCmd())
m.statusMsg = fmt.Sprintf("PR !%s deleted for #%d", ev.PR, issNum)
}
return m, cmds
}
func handleIssueEvent(m Model, ev webhook.IssueEvent, cmds []tea.Cmd) (Model, []tea.Cmd) {
switch ev.Action {
case "opened", "reopened":
_, _ = webhook.ExtractNonce(ev.Body)
cmds = append(cmds, loadCmd())
m.statusMsg = fmt.Sprintf("issue #%d opened", ev.IssueNumber)
case "edited":
cmds = append(cmds, syncIssueSpecPlanCmd(ev.IssueNumber, ev.Body))
m.statusMsg = fmt.Sprintf("issue #%d edited, syncing spec/plan", ev.IssueNumber)
}
return m, cmds
}
func handleIssueCommentEvent(m Model, ev webhook.IssueCommentEvent, cmds []tea.Cmd) (Model, []tea.Cmd) {
if !webhook.IsReviewMarker(ev.CommentBody) {
return m, cmds
}
text := strings.TrimSpace(reviewMarkerStripRe.ReplaceAllString(ev.CommentBody, ""))
if text == "" {
return m, cmds
}
var target *issue.Issue
if ev.PRNumber != 0 {
prStr := strconv.Itoa(ev.PRNumber)
for i := range m.issues {
if m.issues[i].PR == prStr {
target = &m.issues[i]
break
}
}
}
if target == nil {
for i := range m.issues {
if m.issues[i].Number == ev.IssueNumber {
target = &m.issues[i]
break
}
}
}
if target == nil {
m.statusMsg = fmt.Sprintf("审查注入失败:未找到 #%d", ev.IssueNumber)
return m, cmds
}
if !InjectFeedback(m.sessionMgr, *target, text) {
m.statusMsg = fmt.Sprintf("审查注入失败(实施 tab 未找到)#%d", target.Number)
} else {
m.statusMsg = fmt.Sprintf("审查反馈已注入 #%d", target.Number)
}
return m, cmds
}
func handlePushEvent(m Model, ev webhook.PushEvent, cmds []tea.Cmd) (Model, []tea.Cmd) {
settings, err := configLoadFn()
if err != nil {
return m, cmds
}
autoBuild := settings.AutoBuildBranch
if autoBuild == "" {
autoBuild = settings.DevBranch
}
if ev.Branch == settings.DevBranch || ev.Branch == autoBuild {
cmds = append(cmds, branchSyncCmdFn())
}
return m, cmds
}
// mergePRStateCmd writes {pr, branch} into the issue's state JSON.
func mergePRStateCmd(issueNumber int, pr, branch string) tea.Cmd {
return func() tea.Msg {
ctx := context.Background()
root, err := workspaceRootFn()
if err != nil {
return nil
}
remote, err := git.DetectRepo(root)
if err != nil {
return nil
}
token, err := auth.ResolveGiteaToken(remote.Host)
if err != nil {
return nil
}
client := gitea.New(remote.Host, token)
_ = issue.MergeStateJSON(ctx, client, remote.Owner, remote.Repo, issueNumber, map[string]any{
"pr": pr,
"branch": branch,
})
return nil
}
}
// conditionalAutoReviewCmd reads state JSON for the issue, then triggers auto-review
// if ShouldAutoReview is true.
func conditionalAutoReviewCmd(issueNumber int, pr string, issues []issue.Issue) tea.Cmd {
return func() tea.Msg {
ctx := context.Background()
root, err := workspaceRootFn()
if err != nil {
return nil
}
remote, err := git.DetectRepo(root)
if err != nil {
return nil
}
token, err := auth.ResolveGiteaToken(remote.Host)
if err != nil {
return nil
}
client := gitea.New(remote.Host, token)
state, err := issue.ReadStateJSON(ctx, client, remote.Owner, remote.Repo, issueNumber)
if err != nil {
state = map[string]any{}
}
settings, _ := configLoadFn()
if !webhook.ShouldAutoReview(state, settings.AutoReview) {
return nil
}
var worktreePath string
for _, iss := range issues {
if iss.Number == issueNumber {
worktreePath = iss.WorktreePath
break
}
}
var worktreeAbs string
if worktreePath != "" {
worktreeAbs = filepath.Join(root, worktreePath)
} else {
worktreeAbs = root
}
return triggerAutoReviewCmd(issueNumber, pr, worktreeAbs)()
}
}
// triggerAutoReviewCmd runs RunReview fire-and-forget and WatchForNewCodexSession to capture reviewSessionId.
func triggerAutoReviewCmd(issueNumber int, pr string, worktreeAbs string) tea.Cmd {
return func() tea.Msg {
logging.Default.Info("webhook", fmt.Sprintf("触发自动审查 #%d PR=%s", issueNumber, pr))
if worktreeAbs == "" {
root, err := workspaceRootFn()
if err != nil {
return autoReviewCapturedMsg{issueNumber: issueNumber}
}
worktreeAbs = root
}
home, err := os.UserHomeDir()
if err != nil {
return autoReviewCapturedMsg{issueNumber: issueNumber}
}
s, _ := configLoadFn()
prompt := cc.ReviewPrompt(s, struct{ PrNumber string }{PrNumber: pr})
now := time.Now()
sessionsDir := cc.CodexSessionsDir(home, now)
ctx := context.Background()
go cc.RunReview(ctx, cc.ReviewOpts{ //nolint:errcheck
WorkspaceRoot: worktreeAbs,
Prompt: prompt,
})
id, _ := webhookWatchCodexSessionFn(ctx, cc.CodexSessionWatchOpts{
SessionsDir: sessionsDir,
Timeout: 5 * time.Minute,
})
return autoReviewCapturedMsg{issueNumber: issueNumber, sessionID: id}
}
}
// autoReviewPersistCmd writes the captured reviewSessionId to state JSON.
func autoReviewPersistCmd(issueNumber int, sessionID string) tea.Cmd {
return func() tea.Msg {
ctx := context.Background()
root, err := workspaceRootFn()
if err != nil {
return autoReviewPersistMsg{issueNumber: issueNumber, err: err}
}
remote, err := git.DetectRepo(root)
if err != nil {
return autoReviewPersistMsg{issueNumber: issueNumber, err: err}
}
token, err := auth.ResolveGiteaToken(remote.Host)
if err != nil {
return autoReviewPersistMsg{issueNumber: issueNumber, err: err}
}
client := gitea.New(remote.Host, token)
err = issue.MergeStateJSON(ctx, client, remote.Owner, remote.Repo, issueNumber, map[string]any{
"reviewSessionId": sessionID,
})
return autoReviewPersistMsg{issueNumber: issueNumber, err: err}
}
}
// getPullRequestFn is injectable for tests.
var getPullRequestFn = func(ctx context.Context, client *gitea.Client, owner, repo string, number int) (*gitea.PullRequest, error) {
return client.GetPullRequest(ctx, owner, repo, number)
}
// mergeStateJSONFn is injectable for tests.
var mergeStateJSONFn = func(ctx context.Context, client *gitea.Client, owner, repo string, number int, extra map[string]any) error {
return issue.MergeStateJSON(ctx, client, owner, repo, number, extra)
}
// giteaRepoCtx bundles the resolved gitea client + repo coordinates.
type giteaRepoCtx struct {
client *gitea.Client
owner string
repo string
}
// resolveGiteaRepoFn is injectable for tests; resolves workspace → remote → token → client.
var resolveGiteaRepoFn = func() (*giteaRepoCtx, error) {
root, err := workspaceRootFn()
if err != nil {
return nil, err
}
remote, err := git.DetectRepo(root)
if err != nil {
return nil, err
}
token, err := auth.ResolveGiteaToken(remote.Host)
if err != nil {
return nil, err
}
return &giteaRepoCtx{
client: gitea.New(remote.Host, token),
owner: remote.Owner,
repo: remote.Repo,
}, nil
}
// prClosedCmd fetches the PR merged status via Gitea and, if merged, writes
// {prMerged: true, prMergedAt: <ts>} into the issue state JSON.
func prClosedCmd(issueNumber int, pr string) tea.Cmd {
return func() tea.Msg {
prIndex, err := strconv.Atoi(pr)
if err != nil || prIndex <= 0 {
return nil
}
ctx := context.Background()
rc, err := resolveGiteaRepoFn()
if err != nil {
return nil
}
prData, err := getPullRequestFn(ctx, rc.client, rc.owner, rc.repo, prIndex)
if err != nil {
return nil
}
if !prData.Merged {
return nil
}
mergedAt := prData.MergedAt
if mergedAt == "" {
mergedAt = time.Now().UTC().Format(time.RFC3339)
}
_ = mergeStateJSONFn(ctx, rc.client, rc.owner, rc.repo, issueNumber, map[string]any{
"prMerged": true,
"prMergedAt": mergedAt,
})
return nil
}
}
// prDeletedCmd clears the pr field in the issue state JSON.
func prDeletedCmd(issueNumber int) tea.Cmd {
return func() tea.Msg {
ctx := context.Background()
rc, err := resolveGiteaRepoFn()
if err != nil {
return nil
}
_ = mergeStateJSONFn(ctx, rc.client, rc.owner, rc.repo, issueNumber, map[string]any{
"pr": "",
})
return nil
}
}
// syncIssueSpecPlanCmd reads spx:spec / spx:plan markers from the issue body
// and merges them into the state JSON. Faithful port of coordinator.ts handleIssueEdited.
func syncIssueSpecPlanCmd(issueNumber int, body string) tea.Cmd {
return func() tea.Msg {
specMatch := specMarkerRe.FindStringSubmatch(body)
planMatch := planMarkerRe.FindStringSubmatch(body)
if specMatch == nil && planMatch == nil {
return nil
}
extra := make(map[string]any)
if specMatch != nil {
extra["specFile"] = specMatch[1]
}
if planMatch != nil {
extra["planFile"] = planMatch[1]
}
ctx := context.Background()
rc, err := resolveGiteaRepoFn()
if err != nil {
return nil
}
_ = mergeStateJSONFn(ctx, rc.client, rc.owner, rc.repo, issueNumber, extra)
return nil
}
}