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.
This commit is contained in:
wizardchen
2026-05-28 20:16:02 +08:00
committed by lyingbug
parent e3525f884b
commit f1a27e0e18
4 changed files with 91 additions and 0 deletions
@@ -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
}
@@ -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).
@@ -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{
+25
View File
@@ -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{}{