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{}{