diff --git a/backend/plugins/domain/admin/repository/flush_log_test.go b/backend/plugins/domain/admin/repository/flush_log_test.go index fd3eeb57..11004154 100644 --- a/backend/plugins/domain/admin/repository/flush_log_test.go +++ b/backend/plugins/domain/admin/repository/flush_log_test.go @@ -6,7 +6,9 @@ package repository_test import ( "context" "errors" + "fmt" "testing" + "time" "github.com/alicebob/miniredis/v2" "github.com/redis/go-redis/v9" @@ -112,3 +114,37 @@ func TestFlushTaskExecutionLogPersistsAndClears(t *testing.T) { err = svc.Get(ctx, repository.TaskExecutionLogRedisKey(taskID), &buf) assert.True(t, errors.Is(err, contracts.ErrCacheMiss), "flush 后缓存应清空, got %v", err) } + +// readFailCache 读取永远报错而写入成功,用于区分「未命中」与「缓存故障」两种语义。 +type readFailCache struct { + writes []string +} + +func (c *readFailCache) Get(context.Context, string, any) error { + return errors.New("cache unavailable") +} + +func (c *readFailCache) Set(_ context.Context, key string, value any, _ time.Duration) error { + c.writes = append(c.writes, fmt.Sprintf("%s=%v", key, value)) + return nil +} + +func (c *readFailCache) Delete(context.Context, string) error { return nil } + +func (c *readFailCache) GetOrSet(context.Context, string, any, time.Duration, func() (any, error)) error { + return errors.New("cache unavailable") +} + +func (c *readFailCache) Invalidate(context.Context, string) error { return nil } + +// TestAppendTaskExecutionLogKeepsBufferOnCacheReadError 回归:缓存读取失败(而非未命中)时 +// 不得把「读不到」当成「没有缓冲」继续写入,否则整段任务日志会被最新一行覆盖丢失。 +func TestAppendTaskExecutionLogKeepsBufferOnCacheReadError(t *testing.T) { + fake := &readFailCache{} + repository.SetCacheService(fake) + defer repository.SetCacheService(nil) + + err := repository.AppendTaskExecutionLog(context.Background(), "append-err-task", "step-2") + assert.Error(t, err, "缓存故障必须上抛,而不是覆盖缓冲") + assert.Empty(t, fake.writes, "读取失败时不得写入,避免覆盖已缓冲日志") +} diff --git a/backend/plugins/domain/admin/repository/task.go b/backend/plugins/domain/admin/repository/task.go index 4ddf64f9..c6e65b0a 100644 --- a/backend/plugins/domain/admin/repository/task.go +++ b/backend/plugins/domain/admin/repository/task.go @@ -125,7 +125,13 @@ func AppendTaskExecutionLog(ctx context.Context, taskID, logLine string) error { key := TaskExecutionLogRedisKey(taskID) var existing string - _ = cacheSvc.Get(ctx, key, &existing) + if err := cacheSvc.Get(ctx, key, &existing); err != nil { + // 只有未命中才代表「尚无缓冲」;其余读取失败若被当作空缓冲继续写入, + // 会用这一行覆盖掉整段已缓冲的任务日志。 + if !errors.Is(err, contracts.ErrCacheMiss) { + return fmt.Errorf("load buffered task execution log: %w", err) + } + } return cacheSvc.Set(ctx, key, existing+line, taskExecutionLogExpiration) }