Цена читающего маршрута: чекпойнт WAL по таймеру и условный запрос
- рядом с воркером свёртки живёт горутина, раз в минуту разбирающая журнал пассивным чекпойнтом; «журнал не разбирается» видно строкой владельцу, а не только по `df`. Признак — пара чисел, а не флаг занятости: тот молчит под удерживаемым читателем (`busy=0` при 6256 страницах и пяти перенесённых), а при занятой блокировке отдаёт `-1` вместо ответа, и `-1 >= -1` читалось бы как «разобрано целиком» - каталог отвечает `304` на `If-None-Match`, не открывая снимок витрины. Метка собрана из всего, от чего зависит ответ: версии витрины (`data_version` с закреплённого соединения плюс поколение — значение локально для соединения и не переживает переоткрытия), горизонта измерения и области действия ресурса. Версия снимается до и после сборки: снятая после пометила бы устаревший снимок свежим номером - предел и дедлайн ответа отложены в задачу Read API точек вместе с измеренной ценой первого запроса; попутно починен флаки-тест чужой задачи, искавший значение точки в сыром буфере записи лога
This commit is contained in:
@@ -0,0 +1,142 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log/slog"
|
||||
"time"
|
||||
|
||||
"git.vakhrushev.me/av/healthlog/internal/store"
|
||||
)
|
||||
|
||||
// checkpointInterval — как часто разбирается журнал WAL.
|
||||
//
|
||||
// Автоматический чекпойнт SQLite остаётся первой линией и срабатывает по концу
|
||||
// записи; этот тик закрывает случай, которого тот не закрывает по построению —
|
||||
// запись прекратилась, а журнал остался неразобранным. Поток пачечный, ночью
|
||||
// телефон молчит часами, поэтому минута против пяти неразличима по эффекту;
|
||||
// минута взята потому, что с ней своевременен признак «журнал не разбирается»,
|
||||
// и потому, что это тот же ритм, что у тика воркера свёртки. Тот же период
|
||||
// берёт Litestream, у которого задача ровно та же.
|
||||
const checkpointInterval = time.Minute
|
||||
|
||||
// walGrowth — во сколько раз обязан вырасти неразобранный журнал, чтобы о нём
|
||||
// сказали второй раз.
|
||||
//
|
||||
// Признак заводится ради состояния, которое САМО НЕ ПРОХОДИТ: вечный читатель
|
||||
// (в Go чаще всего — незакрытый `sql.Rows`) держит снимок до конца жизни
|
||||
// процесса. Строка на каждый тик дала бы 1440 одинаковых `WARN` в сутки, и
|
||||
// владелец перестал бы их читать раньше, чем кончится диск. Поэтому вторая
|
||||
// строка пишется, только когда стало вдвое хуже.
|
||||
const walGrowth = 2
|
||||
|
||||
// keepWAL разбирает журнал WAL, пока сервис работает.
|
||||
//
|
||||
// Живёт в бинаре, а не в хранилище, и это осознанная асимметрия с воркером
|
||||
// свёртки: у воркера есть доменный исход (доставка свёрнута), а здесь только
|
||||
// жизненный цикл процесса и строка владельцу. Чтобы завести цикл в `store`,
|
||||
// пришлось бы внести туда логгер — первый в пакете, который сегодня не логирует
|
||||
// вовсе и все исходы отдаёт возвратом. Интерпретация чисел при этом осталась в
|
||||
// хранилище (`Checkpoint.Stuck`): семантика тройки `busy/log/checkpointed`
|
||||
// принадлежит SQLite, а не тому, кто её печатает.
|
||||
//
|
||||
// Период параметром, а не константой внутри: тот же шов, что `Worker.Pass` у
|
||||
// свёртки, и по той же причине — иначе проверка «цикл переживает отказ» ждала
|
||||
// бы по минуте на тик. Конфигурируемостью это не является: вызов один, и он
|
||||
// называет константу.
|
||||
//
|
||||
// Контекст один, и работа идёт на нём же — в отличие от свёртки, которая
|
||||
// сворачивает на отвязанном. Прерванный чекпойнт ничего не теряет: перенос
|
||||
// страниц идемпотентен, исхода разбора он не пишет, а следующий старт возьмёт
|
||||
// журнал с того же места. Зато остановка не ждёт переноса полусотни мегабайт в
|
||||
// бюджете, который делится с приёмом и воркером.
|
||||
func keepWAL(ctx context.Context, st *store.Store, log *slog.Logger, every time.Duration) {
|
||||
log = log.With("capability", "wal")
|
||||
ticker := time.NewTicker(every)
|
||||
defer ticker.Stop()
|
||||
|
||||
var watch walWatch
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
}
|
||||
|
||||
ck, err := st.CheckpointWAL(ctx)
|
||||
if err != nil {
|
||||
if errors.Is(err, context.Canceled) {
|
||||
// Штатная остановка не отказ: чекпойнт прерван ею же. ERROR о
|
||||
// ней обесценил бы уровень, по которому вмешиваются, — и делал
|
||||
// бы это на каждом `task restart`.
|
||||
//
|
||||
// Различаем по САМОЙ ошибке, а не по `ctx.Err()`: настоящий
|
||||
// отказ базы, случившийся в тот же тик, что и сигнал остановки,
|
||||
// иначе подавлялся бы как штатный — то есть терялся бы ровно
|
||||
// тогда, когда владелец смотрит в логи.
|
||||
return
|
||||
}
|
||||
// Отказ не прекращает цикл: обслуживание, умершее от временного
|
||||
// отказа базы, молча перестало бы разбирать журнал до конца жизни
|
||||
// процесса — а видно это было бы только по свободному месту.
|
||||
log.ErrorContext(ctx, "wal checkpoint failed", "error", err)
|
||||
continue
|
||||
}
|
||||
|
||||
switch watch.see(ck) {
|
||||
case walStuck:
|
||||
// Адресат — владелец, событие «может стать проблемой»: журнал
|
||||
// растёт, и лечится это не кодом. Значений из данных в записи нет —
|
||||
// только счётчики страниц.
|
||||
log.WarnContext(ctx, "wal checkpoint did not advance",
|
||||
"log_pages", ck.Log,
|
||||
"checkpointed_pages", ck.Checkpointed)
|
||||
case walRecovered:
|
||||
// Возврат к норме — событие, и сказать о нём надо: молчание иначе
|
||||
// неотличимо от «сервис перестал проверять».
|
||||
log.InfoContext(ctx, "wal checkpoint caught up", "log_pages", ck.Log)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// walSay — что сказать владельцу по исходу очередного чекпойнта.
|
||||
type walSay int
|
||||
|
||||
const (
|
||||
walSilent walSay = iota
|
||||
walStuck
|
||||
walRecovered
|
||||
)
|
||||
|
||||
// walWatch решает, когда о неразобранном журнале говорить. Отдельно от цикла,
|
||||
// потому что это единственная его часть, у которой есть исход: решение зависит
|
||||
// от предыдущих тиков, а проверять его ожиданием минут нельзя.
|
||||
type walWatch struct {
|
||||
// warnedAt — размер журнала, о котором уже сказано. Ноль означает
|
||||
// «состояние нормальное». Свойство разговора с владельцем, а не базы,
|
||||
// поэтому живёт здесь, а не в хранилище.
|
||||
warnedAt int
|
||||
}
|
||||
|
||||
func (w *walWatch) see(ck store.Checkpoint) walSay {
|
||||
switch {
|
||||
case !ck.Known():
|
||||
// Исход не измерен (чекпойнт не взял блокировку). Молчим и НЕ трогаем
|
||||
// накопленное: иначе занятый тик посреди беды прочитался бы как
|
||||
// выздоровление, сбросил бы подавитель и вернул те самые 1440 строк в
|
||||
// сутки, против которых он заведён.
|
||||
return walSilent
|
||||
case ck.Stuck() && (w.warnedAt == 0 || ck.Log >= w.warnedAt*walGrowth):
|
||||
w.warnedAt = ck.Log
|
||||
return walStuck
|
||||
case ck.Complete() && w.warnedAt != 0:
|
||||
// Именно `Complete`, а не «порог перестал срабатывать»: журнал, упавший
|
||||
// ниже порога, но так и не перенесённый, — это всё ещё удерживаемый
|
||||
// снимок. Строка «догнали» при нуле перенесённых страниц утверждала бы
|
||||
// то, чего никто не проверял.
|
||||
w.warnedAt = 0
|
||||
return walRecovered
|
||||
default:
|
||||
return walSilent
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,170 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.vakhrushev.me/av/healthlog/internal/store"
|
||||
)
|
||||
|
||||
// Решение «сказать ли владельцу» проверяется таблицей, а не ожиданием минут:
|
||||
// состояние копится по тикам, и без отдельной точки его пришлось бы проверять
|
||||
// прогоном цикла.
|
||||
func TestКогдаГоворитьОНеразобранномЖурнале(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const over = 100000 // заведомо больше порога, выраженного в страницах
|
||||
stuck := store.Checkpoint{Log: over, Checkpointed: 0, PageSize: 4096}
|
||||
worse := store.Checkpoint{Log: over * 4, Checkpointed: 0, PageSize: 4096}
|
||||
slightlyWorse := store.Checkpoint{Log: over + 1, Checkpointed: 0, PageSize: 4096}
|
||||
fine := store.Checkpoint{Log: 12, Checkpointed: 12, PageSize: 4096}
|
||||
// Занятый чекпойнт: исход не измерен, `-1` вместо чисел.
|
||||
unknown := store.Checkpoint{Busy: true, Log: -1, Checkpointed: -1, PageSize: 4096}
|
||||
|
||||
var w walWatch
|
||||
cases := []struct {
|
||||
name string
|
||||
in store.Checkpoint
|
||||
want walSay
|
||||
}{
|
||||
{"первый застрявший чекпойнт", stuck, walStuck},
|
||||
{"то же состояние — молчим", stuck, walSilent},
|
||||
{"чуть хуже — всё ещё молчим", slightlyWorse, walSilent},
|
||||
{"занятый тик посреди беды молчит", unknown, walSilent},
|
||||
{"и не сбрасывает накопленное", stuck, walSilent},
|
||||
{"стало заметно хуже", worse, walStuck},
|
||||
{"разобрался — говорим о возврате", fine, walRecovered},
|
||||
{"норма держится — молчим", fine, walSilent},
|
||||
{"застрял снова", stuck, walStuck},
|
||||
}
|
||||
for _, c := range cases {
|
||||
if got := w.see(c.in); got != c.want {
|
||||
t.Errorf("%s: сказано %v, ждали %v", c.name, got, c.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Цикл обязан пережить отказ базы: обслуживание, умершее от временного отказа,
|
||||
// молча перестало бы разбирать журнал до конца жизни процесса.
|
||||
func TestЦиклЧекпойнтаПереживаетОтказ(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
st, err := store.Open(filepath.Join(t.TempDir(), "healthlog.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("открытие базы: %v", err)
|
||||
}
|
||||
// Закрытая база — самый простой источник устойчивого отказа чекпойнта.
|
||||
if err := st.Close(); err != nil {
|
||||
t.Fatalf("закрытие базы: %v", err)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
seen := &lines{}
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
keepWAL(ctx, st, slog.New(seen), time.Millisecond)
|
||||
}()
|
||||
|
||||
// Даём циклу натолкнуться на отказ много раз подряд.
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
select {
|
||||
case <-done:
|
||||
t.Fatal("цикл вышел сам, не дождавшись отмены")
|
||||
default:
|
||||
}
|
||||
// Отказ обязан быть виден: молча не разбирающийся журнал обнаруживается
|
||||
// только по свободному месту.
|
||||
if !seen.has("wal checkpoint failed") {
|
||||
t.Error("отказ чекпойнта не оставил записи владельцу")
|
||||
}
|
||||
|
||||
cancel()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("цикл не вышел по отмене")
|
||||
}
|
||||
}
|
||||
|
||||
// Отмена — единственный законный повод выйти, и выйти надо сразу: горутина
|
||||
// ждётся в общем бюджете остановки вместе с воркером свёртки.
|
||||
func TestЦиклЧекпойнтаВыходитПоОтмене(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
st, err := store.Open(filepath.Join(t.TempDir(), "healthlog.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("открытие базы: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = st.Close() })
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
keepWAL(ctx, st, slog.New(slog.DiscardHandler), time.Millisecond)
|
||||
}()
|
||||
|
||||
cancel()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("цикл не вышел по отмене")
|
||||
}
|
||||
}
|
||||
|
||||
// Ветка «фоновые горутины не уложились в бюджет» — последняя защита инварианта
|
||||
// «доставка либо свёрнута целиком, либо остаётся pending». Прогоном сервиса её
|
||||
// не проверить: бюджет тридцать секунд, а заставить воркер зависнуть нечем.
|
||||
func TestОжиданиеФоновыхГорутин(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
closed := make(chan struct{})
|
||||
close(closed)
|
||||
if !waitBackground(context.Background(), closed, slog.New(slog.DiscardHandler)) {
|
||||
t.Error("вышедшие горутины не дождались")
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
seen := &lines{}
|
||||
if waitBackground(ctx, make(chan struct{}), slog.New(seen)) {
|
||||
t.Error("зависшие горутины объявлены вышедшими — база закрылась бы из-под них")
|
||||
}
|
||||
if !seen.has("shutdown budget exceeded") {
|
||||
t.Error("превышение бюджета осталось без строки владельцу")
|
||||
}
|
||||
}
|
||||
|
||||
// lines — slog.Handler, копящий сообщения: проверяется факт записи, не данные.
|
||||
type lines struct {
|
||||
mu sync.Mutex
|
||||
msg []string
|
||||
}
|
||||
|
||||
func (l *lines) Enabled(context.Context, slog.Level) bool { return true }
|
||||
|
||||
func (l *lines) Handle(_ context.Context, rec slog.Record) error {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
l.msg = append(l.msg, rec.Message)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (l *lines) WithAttrs([]slog.Attr) slog.Handler { return l }
|
||||
func (l *lines) WithGroup(string) slog.Handler { return l }
|
||||
|
||||
func (l *lines) has(msg string) bool {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
for _, m := range l.msg {
|
||||
if m == msg {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
+53
-15
@@ -9,6 +9,7 @@ import (
|
||||
"net"
|
||||
"net/http"
|
||||
"os/signal"
|
||||
"sync"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
@@ -46,6 +47,31 @@ func runServe(args []string) error {
|
||||
return serve(ctx, cfg, logging.New(cfg.Log.Level, cfg.Log.Format), nil)
|
||||
}
|
||||
|
||||
// waitBackground ждёт выхода фоновых горутин и говорит, дождался ли.
|
||||
//
|
||||
// Отдельной функцией потому, что это единственная ветка остановки, у которой
|
||||
// есть исход, и проверить её прогоном сервиса нельзя: бюджет — тридцать секунд,
|
||||
// а заставить воркер зависнуть по требованию нечем.
|
||||
//
|
||||
// Не дождались — база НЕ закрывается: её транзакцию свернёт выход процесса, и
|
||||
// доставка останется `pending`, то есть будет подобрана следующим стартом.
|
||||
// Закрытая из-под воркера, она дала бы ERROR по доставке, с которой всё в
|
||||
// порядке.
|
||||
//
|
||||
// Этап в записи называется общим именем, а не воркером свёртки: ждём мы двоих,
|
||||
// и назвать виновным одного из них значило бы угадать. Чекпойнт при этом
|
||||
// выходит по отмене немедленно, так что практически это всё тот же воркер, — но
|
||||
// лог не должен утверждать того, чего не проверял.
|
||||
func waitBackground(shutdownCtx context.Context, done <-chan struct{}, log *slog.Logger) bool {
|
||||
select {
|
||||
case <-done:
|
||||
return true
|
||||
case <-shutdownCtx.Done():
|
||||
log.Warn("shutdown budget exceeded", "stage", "background")
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// serve поднимает сервис и ведёт его до отмены контекста.
|
||||
//
|
||||
// Контекст параметром, а не подпиской на сигнал внутри: иначе весь жизненный
|
||||
@@ -119,23 +145,41 @@ func serve(ctx context.Context, cfg *config.Config, log *slog.Logger, ready func
|
||||
return fmt.Errorf("listen %q: %w", cfg.Server.Addr, err)
|
||||
}
|
||||
|
||||
workerCtx, stopWorker := context.WithCancel(context.Background())
|
||||
defer stopWorker()
|
||||
workerDone := make(chan struct{})
|
||||
go func() {
|
||||
defer close(workerDone)
|
||||
// Обе фоновые горутины живут на одном контексте и ждутся вместе. Вместе —
|
||||
// потому что база закрывается ПОСЛЕ выхода обеих: закрытая из-под воркера,
|
||||
// она даёт ERROR по доставке, с которой всё в порядке, а из-под чекпойнта —
|
||||
// отказ обслуживания на ровном месте.
|
||||
bgCtx, stopBackground := context.WithCancel(context.Background())
|
||||
defer stopBackground()
|
||||
|
||||
var bg sync.WaitGroup
|
||||
bg.Go(func() {
|
||||
// Первый проход воркера и есть подбор неразобранного при старте:
|
||||
// отдельного кода для него нет намеренно.
|
||||
worker.Run(workerCtx)
|
||||
worker.Run(bgCtx)
|
||||
})
|
||||
bg.Go(func() {
|
||||
keepWAL(bgCtx, st, log, checkpointInterval)
|
||||
})
|
||||
backgroundDone := make(chan struct{})
|
||||
go func() {
|
||||
bg.Wait()
|
||||
close(backgroundDone)
|
||||
}()
|
||||
|
||||
errCh := make(chan error, 1)
|
||||
go func() {
|
||||
// Параметры обслуживания журнала — в той же строке, а не отдельной:
|
||||
// горутина, которую забыли запустить, иначе неотличима от здоровой
|
||||
// ровно до того дня, когда журнал упрётся в диск. Ноль новых строк, обе
|
||||
// константы проверяемы глазами.
|
||||
log.Info("server started",
|
||||
"addr", ln.Addr().String(),
|
||||
"db_path", cfg.Storage.DBPath,
|
||||
"archive_dir", arch.Root(),
|
||||
"max_body_mb", cfg.Ingest.MaxBodyMB)
|
||||
"max_body_mb", cfg.Ingest.MaxBodyMB,
|
||||
"wal_checkpoint_sec", int64(checkpointInterval.Seconds()),
|
||||
"wal_limit_mb", store.JournalSizeLimitMB)
|
||||
|
||||
if err := srv.Serve(ln); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||
errCh <- fmt.Errorf("serve: %w", err)
|
||||
@@ -172,15 +216,9 @@ func serve(ctx context.Context, cfg *config.Config, log *slog.Logger, ready func
|
||||
}
|
||||
}
|
||||
|
||||
stopWorker()
|
||||
select {
|
||||
case <-workerDone:
|
||||
stopBackground()
|
||||
if waitBackground(shutdownCtx, backgroundDone, log) {
|
||||
closeStore()
|
||||
case <-shutdownCtx.Done():
|
||||
// Воркер не вышел в бюджет. База не закрывается: её транзакцию свернёт
|
||||
// выход процесса, и доставка останется `pending` — то есть будет
|
||||
// подобрана следующим стартом.
|
||||
log.Warn("shutdown budget exceeded", "stage", "fold-worker")
|
||||
}
|
||||
return serveErr
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user