130 lines
2.5 KiB
Go
130 lines
2.5 KiB
Go
package store
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"errors"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// ErrQueueClosed 表示写入队列已关闭。
|
|
var ErrQueueClosed = errors.New("store: write queue closed")
|
|
|
|
// ErrBusy 表示写库基础设施失败(可映射为协议 busy;/readyz 应失败)。
|
|
var ErrBusy = errors.New("store: busy")
|
|
|
|
// WriteFunc 在单个写事务中执行的操作。
|
|
type WriteFunc func(tx *sql.Tx) error
|
|
|
|
// Queue 写入队列:提交一个写操作并拿到结果。
|
|
//
|
|
// P1 仍为一操作一事务;P2 换成合并提交(最多 256 / 2ms + SAVEPOINT)。
|
|
type Queue struct {
|
|
db *sql.DB
|
|
|
|
mu sync.Mutex
|
|
closed bool
|
|
inflight int
|
|
ready bool
|
|
lastWriteErr error
|
|
}
|
|
|
|
// NewQueue 创建简单写入队列(一操作一事务)。
|
|
func NewQueue(db *sql.DB) *Queue {
|
|
return &Queue{db: db, ready: true}
|
|
}
|
|
|
|
// Do 提交写操作并等待提交结果。
|
|
func (q *Queue) Do(ctx context.Context, fn WriteFunc) error {
|
|
if fn == nil {
|
|
return errors.New("store: nil write func")
|
|
}
|
|
q.mu.Lock()
|
|
if q.closed {
|
|
q.mu.Unlock()
|
|
return ErrQueueClosed
|
|
}
|
|
q.inflight++
|
|
q.mu.Unlock()
|
|
|
|
defer func() {
|
|
q.mu.Lock()
|
|
q.inflight--
|
|
q.mu.Unlock()
|
|
}()
|
|
|
|
if err := ctx.Err(); err != nil {
|
|
return err
|
|
}
|
|
tx, err := q.db.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
q.markBusy(err)
|
|
return errors.Join(ErrBusy, err)
|
|
}
|
|
if err := fn(tx); err != nil {
|
|
_ = tx.Rollback()
|
|
return err
|
|
}
|
|
if err := tx.Commit(); err != nil {
|
|
q.markBusy(err)
|
|
return errors.Join(ErrBusy, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (q *Queue) markBusy(err error) {
|
|
q.mu.Lock()
|
|
defer q.mu.Unlock()
|
|
q.ready = false
|
|
q.lastWriteErr = err
|
|
}
|
|
|
|
// IsReady 写库是否仍可用(写失败后为 false)。
|
|
func (q *Queue) IsReady() bool {
|
|
q.mu.Lock()
|
|
defer q.mu.Unlock()
|
|
return q.ready
|
|
}
|
|
|
|
// LastWriteError 返回最近一次基础设施写失败。
|
|
func (q *Queue) LastWriteError() error {
|
|
q.mu.Lock()
|
|
defer q.mu.Unlock()
|
|
return q.lastWriteErr
|
|
}
|
|
|
|
// Len 返回进行中的写操作数。
|
|
func (q *Queue) Len() int {
|
|
q.mu.Lock()
|
|
defer q.mu.Unlock()
|
|
return q.inflight
|
|
}
|
|
|
|
// Drain 等待已进行中的写操作完成,或 ctx 取消。
|
|
func (q *Queue) Drain(ctx context.Context) error {
|
|
ticker := time.NewTicker(5 * time.Millisecond)
|
|
defer ticker.Stop()
|
|
for {
|
|
q.mu.Lock()
|
|
n := q.inflight
|
|
q.mu.Unlock()
|
|
if n == 0 {
|
|
return nil
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case <-ticker.C:
|
|
}
|
|
}
|
|
}
|
|
|
|
// Close 关闭队列,之后 Do 返回 ErrQueueClosed。
|
|
func (q *Queue) Close() error {
|
|
q.mu.Lock()
|
|
defer q.mu.Unlock()
|
|
q.closed = true
|
|
return nil
|
|
}
|