diff --git a/docs/DEVIATIONS.md b/docs/DEVIATIONS.md index cc28178..58f2504 100644 --- a/docs/DEVIATIONS.md +++ b/docs/DEVIATIONS.md @@ -1728,3 +1728,13 @@ issue #3 未关闭,`feat/fix-3-downlink-deadlock` 未合入 `main`。下面是 - 原因:终态回执、超限后的后续 pending、撤回腾出窗口后都依赖再唤醒,否则空闲在线端收不到。 - 备选方案:在同一次 `PushPending` 里循环扫完整 pending 表(否决,指令要求合并唤醒下一轮)。 - 影响:仅 `internal/app/message/`;相关单测见 `review_r3_65_test.go`。 + +### 复审修复 R3-03 + +1. **写队列满时 Close 不再与 enqueue 死锁** + - 日期:2026-09-30 + - 原条款:issue #67;停机需 Drain/Close 在超时内返回;已入队任务仍由写 goroutine 处理完。 + - 实际做法:Queue 增加 closing 信号。enqueue 在持 sendMu.RLock 时 select 同时等待 q.ch、closing 与 ctx.Done(),通道满时不再无限阻塞。Close 先 Swap(closed) 并 close(closing) 唤醒在途发送方释放读锁,再在写锁内 close(q.ch),最后等 loop 退出。不向已关闭 channel 发送。 + - 原因:通道满时并发 Do 占着读锁堵在发送上,Close 的写锁拿不到,停机卡死;未提交的在途写也会丢。 + - 备选方案:入队改为非阻塞,满则立即 ErrBusy(否决:改变背压语义,正常高峰会误伤提交)。 + - 影响:仅 internal/store/queue.go 与测试;不改迁移编号。 diff --git a/internal/store/queue.go b/internal/store/queue.go index 99ffb1c..20b52e0 100644 --- a/internal/store/queue.go +++ b/internal/store/queue.go @@ -36,10 +36,11 @@ type writeJob struct { type Queue struct { db *sql.DB - ch chan writeJob - done chan struct{} - closed atomic.Bool - sendMu sync.RWMutex + ch chan writeJob + done chan struct{} + closing chan struct{} // Close 时关闭,唤醒持读锁阻塞在发送上的 enqueue + closed atomic.Bool + sendMu sync.RWMutex mu sync.Mutex ready bool @@ -52,11 +53,16 @@ type Queue struct { // NewQueue 创建合并写入队列并启动写 goroutine。 func NewQueue(db *sql.DB) *Queue { + return newQueue(db, queueBuffSize) +} + +func newQueue(db *sql.DB, buffSize int) *Queue { q := &Queue{ - db: db, - ch: make(chan writeJob, queueBuffSize), - done: make(chan struct{}), - ready: true, + db: db, + ch: make(chan writeJob, buffSize), + done: make(chan struct{}), + closing: make(chan struct{}), + ready: true, } go q.loop() return q @@ -115,10 +121,15 @@ func (q *Queue) enqueue(job writeJob) error { return ErrQueueClosed } q.addPending(1) + // 通道满时不得只堵在发送上持有读锁:Close 需要写锁关闭 q.ch。 select { case q.ch <- job: q.sendMu.RUnlock() return nil + case <-q.closing: + q.addPending(-1) + q.sendMu.RUnlock() + return ErrQueueClosed case <-job.ctx.Done(): q.addPending(-1) q.sendMu.RUnlock() @@ -376,11 +387,13 @@ func (q *Queue) Drain(ctx context.Context) error { } // Close 关闭队列:不再接受新任务,并等待写 goroutine 处理完已入队任务后退出。 -// 在写锁内关闭数据通道,避免并发 Do 向已关闭 channel 发送而 panic。 +// 先关闭 closing 唤醒因通道满而阻塞的发送方并释放读锁,再在无发送者时关闭数据通道, +// 避免向已关闭 channel 发送而 panic,也避免与持读锁的 enqueue 死锁。 func (q *Queue) Close() error { if q.closed.Swap(true) { return nil } + close(q.closing) q.sendMu.Lock() close(q.ch) q.sendMu.Unlock() diff --git a/internal/store/queue_test.go b/internal/store/queue_test.go index d51cb26..7f7dede 100644 --- a/internal/store/queue_test.go +++ b/internal/store/queue_test.go @@ -289,3 +289,137 @@ ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = excluded.upd } } } + +// openSmallQueue 用小缓冲队列替换默认队列,便于测满通道时的 Close/Drain。 +func openSmallQueue(t *testing.T, buf int) *DB { + t.Helper() + dir := t.TempDir() + db, err := Open(dir, "FULL") + if err != nil { + t.Fatal(err) + } + if err := db.Queue.Close(); err != nil { + t.Fatal(err) + } + db.Queue = newQueue(db.Write, buf) + return db +} + +func TestQueueCloseUnblocksFullChannel(t *testing.T) { + t.Parallel() + const buf = 4 + db := openSmallQueue(t, buf) + defer func() { _ = db.Close() }() + q := db.Queue + + hold := make(chan struct{}) + blockerStarted := make(chan struct{}) + blockerErr := make(chan error, 1) + go func() { + blockerErr <- q.Do(context.Background(), func(tx *sql.Tx) error { + close(blockerStarted) + <-hold + return nil + }) + }() + <-blockerStarted + + ctx := context.Background() + var fillWG sync.WaitGroup + for i := 0; i < buf; i++ { + fillWG.Add(1) + go func() { + defer fillWG.Done() + _ = q.Do(ctx, func(tx *sql.Tx) error { return nil }) + }() + } + deadline := time.Now().Add(2 * time.Second) + for q.Len() < buf+1 && time.Now().Before(deadline) { + time.Sleep(2 * time.Millisecond) + } + if q.Len() < buf+1 { + t.Fatalf("channel not full: len=%d", q.Len()) + } + + blockedErr := make(chan error, 1) + go func() { + blockedErr <- q.Do(ctx, func(tx *sql.Tx) error { return nil }) + }() + // 等额外 Do 堵在 enqueue(pending 超过通道容量+正在执行的一条)。 + deadline = time.Now().Add(2 * time.Second) + for q.Len() < buf+2 && time.Now().Before(deadline) { + time.Sleep(2 * time.Millisecond) + } + + closeDone := make(chan error, 1) + go func() { closeDone <- q.Close() }() + + select { + case err := <-blockedErr: + if !errors.Is(err, ErrQueueClosed) { + t.Fatalf("blocked Do: %v", err) + } + case <-time.After(2 * time.Second): + t.Fatal("Close did not unblock full-channel enqueue") + } + + close(hold) + select { + case err := <-closeDone: + if err != nil { + t.Fatal(err) + } + case <-time.After(2 * time.Second): + t.Fatal("Close hung after writer released") + } + <-blockerErr + fillWG.Wait() +} + +func TestQueueDrainTimeoutWhileWriterBlocked(t *testing.T) { + t.Parallel() + const buf = 4 + db := openSmallQueue(t, buf) + defer func() { _ = db.Close() }() + q := db.Queue + + hold := make(chan struct{}) + started := make(chan struct{}) + go func() { + _ = q.Do(context.Background(), func(tx *sql.Tx) error { + close(started) + <-hold + return nil + }) + }() + <-started + + ctx := context.Background() + var wg sync.WaitGroup + for i := 0; i < buf; i++ { + wg.Add(1) + go func() { + defer wg.Done() + _ = q.Do(ctx, func(tx *sql.Tx) error { return nil }) + }() + } + deadline := time.Now().Add(2 * time.Second) + for q.Len() < buf+1 && time.Now().Before(deadline) { + time.Sleep(2 * time.Millisecond) + } + + drainCtx, cancel := context.WithTimeout(context.Background(), 80*time.Millisecond) + defer cancel() + start := time.Now() + err := q.Drain(drainCtx) + elapsed := time.Since(start) + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("Drain err=%v want deadline", err) + } + if elapsed > 500*time.Millisecond { + t.Fatalf("Drain took %s, should return on timeout", elapsed) + } + + close(hold) + wg.Wait() +}