From f1a27e0e18eaec88e35e8246ce27fc403ce5e73d Mon Sep 17 00:00:00 2001 From: wizardchen Date: Thu, 28 May 2026 20:09:49 +0800 Subject: [PATCH] feat: implement CancelOpenSpansByName method and related tests This commit introduces the CancelOpenSpansByName method in the KnowledgeSpanRepository, allowing for the cancellation of open spans by their name for a specific knowledge ID and attempt. This functionality is crucial for managing spans during retries or server restarts, preventing duplicate entries in the trace tree. Additionally, a new test case, TestKnowledgeSpanRepo_CancelOpenSpansByName, has been added to ensure the correct behavior of this method, verifying that only the intended spans are cancelled while others remain unaffected. This enhancement improves the robustness of span management in the application. --- .../repository/knowledge_span_repo.go | 29 +++++++++++++++++++ .../repository/knowledge_span_repo_test.go | 29 +++++++++++++++++++ .../service/knowledge_span_tracker.go | 8 +++++ internal/container/container.go | 25 ++++++++++++++++ 4 files changed, 91 insertions(+) diff --git a/internal/application/repository/knowledge_span_repo.go b/internal/application/repository/knowledge_span_repo.go index ec4e3623a..1defb0f78 100644 --- a/internal/application/repository/knowledge_span_repo.go +++ b/internal/application/repository/knowledge_span_repo.go @@ -38,6 +38,11 @@ type KnowledgeSpanRepository interface { // running — a tree walk that stops at terminal parents would miss // those orphan leaves. CancelAllOpenSpans(ctx context.Context, knowledgeID string, attempt int, errorCode, reason string) (int64, error) + // CancelOpenSpansByName flips pending/running rows with the given span + // name for (knowledgeID, attempt). Used before re-opening a subspan + // after asynq retry or server restart so the trace tree does not + // accumulate duplicate postprocess.summary / question rows. + CancelOpenSpansByName(ctx context.Context, knowledgeID string, attempt int, name, errorCode, reason string) (int64, error) } type knowledgeSpanRepository struct { @@ -228,3 +233,27 @@ func (r *knowledgeSpanRepository) CancelAllOpenSpans( } return res.RowsAffected, nil } + +func (r *knowledgeSpanRepository) CancelOpenSpansByName( + ctx context.Context, knowledgeID string, attempt int, name, errorCode, reason string, +) (int64, error) { + if knowledgeID == "" || attempt <= 0 || name == "" { + return 0, nil + } + now := time.Now() + res := r.db.WithContext(ctx).Model(&types.KnowledgeProcessingSpan{}). + Where("knowledge_id = ? AND attempt = ? AND name = ? AND status IN ?", + knowledgeID, attempt, name, + []string{types.SpanStatusPending, types.SpanStatusRunning}). + Updates(map[string]any{ + "status": types.SpanStatusCancelled, + "error_code": errorCode, + "error_message": reason, + "finished_at": now, + "updated_at": now, + }) + if res.Error != nil { + return 0, res.Error + } + return res.RowsAffected, nil +} diff --git a/internal/application/repository/knowledge_span_repo_test.go b/internal/application/repository/knowledge_span_repo_test.go index 482996e87..3b0675672 100644 --- a/internal/application/repository/knowledge_span_repo_test.go +++ b/internal/application/repository/knowledge_span_repo_test.go @@ -153,6 +153,35 @@ func TestKnowledgeSpanRepo_CancelDescendants(t *testing.T) { assert.Equal(t, types.SpanStatusDone, statusBy["image0"], "terminal states must not be touched") } +func TestKnowledgeSpanRepo_CancelOpenSpansByName(t *testing.T) { + repo, _ := setupSpanTestRepo(t) + ctx := context.Background() + kid := "kid-supersede" + now := time.Now() + + for _, r := range []*types.KnowledgeProcessingSpan{ + {KnowledgeID: kid, Attempt: 1, SpanID: "sum-old", Name: "postprocess.summary", Kind: types.SpanKindSubSpan, Status: types.SpanStatusRunning, StartedAt: &now}, + {KnowledgeID: kid, Attempt: 1, SpanID: "sum-done", Name: "postprocess.summary", Kind: types.SpanKindSubSpan, Status: types.SpanStatusDone, StartedAt: &now}, + {KnowledgeID: kid, Attempt: 1, SpanID: "q-old", Name: "postprocess.question", Kind: types.SpanKindSubSpan, Status: types.SpanStatusRunning, StartedAt: &now}, + } { + require.NoError(t, repo.Upsert(ctx, r)) + } + + affected, err := repo.CancelOpenSpansByName(ctx, kid, 1, "postprocess.summary", "TASK_SUPERSEDED", "retry") + require.NoError(t, err) + assert.Equal(t, int64(1), affected) + + rows, err := repo.ListByAttempt(ctx, kid, 1) + require.NoError(t, err) + statusBy := map[string]string{} + for _, r := range rows { + statusBy[r.SpanID] = r.Status + } + assert.Equal(t, types.SpanStatusCancelled, statusBy["sum-old"]) + assert.Equal(t, types.SpanStatusDone, statusBy["sum-done"]) + assert.Equal(t, types.SpanStatusRunning, statusBy["q-old"]) +} + // TestKnowledgeSpanRepo_ListAttemptIsolation guarantees that different // attempts of the same knowledge stay queryable independently — the // foundation for the "show history" UI navigation (?attempt=N). diff --git a/internal/application/service/knowledge_span_tracker.go b/internal/application/service/knowledge_span_tracker.go index c62974836..501a068e7 100644 --- a/internal/application/service/knowledge_span_tracker.go +++ b/internal/application/service/knowledge_span_tracker.go @@ -369,6 +369,14 @@ func (t *spanTracker) BeginSubSpan(ctx context.Context, parent *Span, name, kind if kind != types.SpanKindGeneration && kind != types.SpanKindSubSpan { kind = types.SpanKindSubSpan } + // Asynq retry / server restart can re-run the same handler while the + // previous invocation's span is still status=running (worker died + // without EndSpan). Cancel same-name open rows so the UI shows one + // logical subspan per (attempt, name) instead of duplicate stripes. + if _, err := t.repo.CancelOpenSpansByName(ctx, parent.KnowledgeID, parent.Attempt, name, + "TASK_SUPERSEDED", "superseded by a new run of the same subtask"); err != nil { + logger.Warnf(ctx, "[SpanTracker] supersede %s before BeginSubSpan failed: %v", name, err) + } now := time.Now() id := newSpanID() row := &types.KnowledgeProcessingSpan{ diff --git a/internal/container/container.go b/internal/container/container.go index c3f583df9..0ea27f119 100644 --- a/internal/container/container.go +++ b/internal/container/container.go @@ -649,6 +649,8 @@ func syncSequences(db *gorm.DB) { // Newer rows are left alone so we don't race a peer instance that's mid-process. func resetPendingTasks(db *gorm.DB) { distributed := os.Getenv("REDIS_ADDR") != "" + ctx := context.Background() + spanRepo := repository.NewKnowledgeSpanRepository(db) knowledgeQuery := db.Model(&types.Knowledge{}). Where("parse_status IN ?", []string{ @@ -672,6 +674,29 @@ func resetPendingTasks(db *gorm.DB) { syncQuery = syncQuery.Where("start_time < ?", staleCutoff) } + // Cancel orphaned trace spans for knowledge rows we are about to mark + // failed. resetPendingTasks does not touch asynq queues; this only + // prevents the UI from showing duplicate running postprocess.* + // subspans when a later retry also opens fresh spans. + var stuckKnowledge []types.Knowledge + if err := knowledgeQuery.Select("id").Find(&stuckKnowledge).Error; err != nil { + logger.Warnf(ctx, "resetPendingTasks: list stuck knowledge failed: %v", err) + } else { + for _, k := range stuckKnowledge { + attempt, err := spanRepo.LatestAttempt(ctx, k.ID) + if err != nil || attempt <= 0 { + continue + } + if n, err := spanRepo.CancelAllOpenSpans(ctx, k.ID, attempt, + "SERVER_RESTART", "task interrupted due to application restart"); err != nil { + logger.Warnf(ctx, "resetPendingTasks: cancel spans for %s failed: %v", k.ID, err) + } else if n > 0 { + logger.Infof(ctx, "resetPendingTasks: cancelled %d open span(s) for knowledge %s attempt %d", + n, k.ID, attempt) + } + } + } + // 1. Reset knowledge parsing tasks (including finalizing rows whose // enrichment subtasks were lost with the process). result := knowledgeQuery.Updates(map[string]interface{}{