- рядом с воркером свёртки живёт горутина, раз в минуту разбирающая журнал пассивным чекпойнтом; «журнал не разбирается» видно строкой владельцу, а не только по `df`. Признак — пара чисел, а не флаг занятости: тот молчит под удерживаемым читателем (`busy=0` при 6256 страницах и пяти перенесённых), а при занятой блокировке отдаёт `-1` вместо ответа, и `-1 >= -1` читалось бы как «разобрано целиком» - каталог отвечает `304` на `If-None-Match`, не открывая снимок витрины. Метка собрана из всего, от чего зависит ответ: версии витрины (`data_version` с закреплённого соединения плюс поколение — значение локально для соединения и не переживает переоткрытия), горизонта измерения и области действия ресурса. Версия снимается до и после сборки: снятая после пометила бы устаревший снимок свежим номером - предел и дедлайн ответа отложены в задачу Read API точек вместе с измеренной ценой первого запроса; попутно починен флаки-тест чужой задачи, искавший значение точки в сыром буфере записи лога
225 lines
10 KiB
Go
225 lines
10 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"log/slog"
|
|
"net"
|
|
"net/http"
|
|
"os/signal"
|
|
"sync"
|
|
"syscall"
|
|
"time"
|
|
|
|
"git.vakhrushev.me/av/healthlog/internal/archive"
|
|
"git.vakhrushev.me/av/healthlog/internal/catalog"
|
|
"git.vakhrushev.me/av/healthlog/internal/config"
|
|
"git.vakhrushev.me/av/healthlog/internal/fold"
|
|
"git.vakhrushev.me/av/healthlog/internal/httpapi"
|
|
"git.vakhrushev.me/av/healthlog/internal/ingest"
|
|
"git.vakhrushev.me/av/healthlog/internal/logging"
|
|
"git.vakhrushev.me/av/healthlog/internal/replay"
|
|
"git.vakhrushev.me/av/healthlog/internal/store"
|
|
)
|
|
|
|
// shutdownTimeout — общий бюджет остановки: сперва дожидаемся активных
|
|
// запросов, затем выхода воркера свёртки. Совпадает со `stop_grace_period`
|
|
// контейнера — за его пределом процесс всё равно убивают.
|
|
const shutdownTimeout = 30 * time.Second
|
|
|
|
func runServe(args []string) error {
|
|
fs := flag.NewFlagSet("serve", flag.ContinueOnError)
|
|
cfgPath := fs.String("config", config.DefaultPath, "путь к config.toml")
|
|
if err := fs.Parse(args); err != nil {
|
|
return fmt.Errorf("parse flags: %w", err)
|
|
}
|
|
|
|
cfg, err := config.Load(*cfgPath)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
|
defer stop()
|
|
|
|
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 поднимает сервис и ведёт его до отмены контекста.
|
|
//
|
|
// Контекст параметром, а не подпиской на сигнал внутри: иначе весь жизненный
|
|
// цикл — порядок остановки, ожидание воркера, судьба несвёрнутой доставки —
|
|
// проверялся бы только посылкой сигнала самому себе, то есть не проверялся бы.
|
|
//
|
|
// ready, если задан, зовётся с ФАКТИЧЕСКИМ адресом прослушивания: при `:0` в
|
|
// конфиге узнать порт больше неоткуда.
|
|
func serve(ctx context.Context, cfg *config.Config, log *slog.Logger, ready func(addr string)) error {
|
|
st, err := store.Open(cfg.Storage.DBPath)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// Закрытие базы — не `defer`: при исчерпании бюджета остановки воркер может
|
|
// ещё сворачивать доставку, и закрытая из-под него база дала бы ERROR по
|
|
// доставке, с которой всё в порядке. Кто закрывает, решает ветка остановки.
|
|
closed := false
|
|
closeStore := func() {
|
|
if !closed {
|
|
closed = true
|
|
_ = st.Close()
|
|
}
|
|
}
|
|
|
|
arch, err := archive.New(cfg.Storage.ArchiveDir)
|
|
if err != nil {
|
|
closeStore()
|
|
return err
|
|
}
|
|
|
|
if len(cfg.Auth.WriteTokens) == 0 {
|
|
log.Warn("write auth disabled", "reason", "auth.write_tokens пуст")
|
|
}
|
|
// Цена у двух контуров разная, и это сказано вслух: открытый приём означает
|
|
// мусор во входе, открытое чтение — выгрузку истории здоровья любому, кто
|
|
// нашёл порт. Пока сервис живёт в доверенной сети, это осознанный выбор;
|
|
// перед выкладкой наружу список обязан быть непуст.
|
|
if len(cfg.Auth.ReadTokens) == 0 {
|
|
log.Warn("read auth disabled", "reason", "auth.read_tokens пуст")
|
|
}
|
|
|
|
// Воркер и приём делят одну свёртку: приём её только будит, сворачивает
|
|
// воркер — и в порядке журнала, чего синхронная свёртка внутри обработчика
|
|
// не давала при конкурентных доставках.
|
|
worker := replay.NewWorker(st, fold.New(arch, st, int64(cfg.Ingest.MaxBodyMB)<<20, log), log)
|
|
|
|
srv := &http.Server{
|
|
Handler: httpapi.New(httpapi.Options{
|
|
Ingest: ingest.New(arch, st, worker.Notify, log),
|
|
Catalog: catalog.New(st, log),
|
|
Log: log,
|
|
WriteTokens: cfg.Auth.WriteTokens,
|
|
ReadTokens: cfg.Auth.ReadTokens,
|
|
MaxBodyMB: cfg.Ingest.MaxBodyMB,
|
|
// Бюджет ответа маршрута приёма: `WriteTimeout` сервера ставится ДО
|
|
// вызова обработчика и потому покрывает чтение тела, обрывая
|
|
// медленную загрузку молча. Длинный бюджет нужен одному маршруту,
|
|
// поэтому и выдаётся ему, а не всему серверу.
|
|
IngestWriteBudget: cfg.Server.ReadTimeout.D() + cfg.Server.WriteTimeout.D(),
|
|
}),
|
|
// ReadTimeout щедрый (большой пакет по мобильной сети), но заголовки
|
|
// обязаны приехать быстро — иначе полуоткрытое соединение держит слот.
|
|
ReadHeaderTimeout: 10 * time.Second,
|
|
ReadTimeout: cfg.Server.ReadTimeout.D(),
|
|
WriteTimeout: cfg.Server.WriteTimeout.D(),
|
|
}
|
|
|
|
ln, err := net.Listen("tcp", cfg.Server.Addr)
|
|
if err != nil {
|
|
closeStore()
|
|
return fmt.Errorf("listen %q: %w", cfg.Server.Addr, err)
|
|
}
|
|
|
|
// Обе фоновые горутины живут на одном контексте и ждутся вместе. Вместе —
|
|
// потому что база закрывается ПОСЛЕ выхода обеих: закрытая из-под воркера,
|
|
// она даёт ERROR по доставке, с которой всё в порядке, а из-под чекпойнта —
|
|
// отказ обслуживания на ровном месте.
|
|
bgCtx, stopBackground := context.WithCancel(context.Background())
|
|
defer stopBackground()
|
|
|
|
var bg sync.WaitGroup
|
|
bg.Go(func() {
|
|
// Первый проход воркера и есть подбор неразобранного при старте:
|
|
// отдельного кода для него нет намеренно.
|
|
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,
|
|
"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)
|
|
}
|
|
}()
|
|
if ready != nil {
|
|
ready(ln.Addr().String())
|
|
}
|
|
|
|
var serveErr error
|
|
select {
|
|
case serveErr = <-errCh:
|
|
// Отказ приёма не отменяет остановки воркера: закрыть базу, не дождавшись
|
|
// его, значит выдернуть её из-под идущей свёртки и получить ERROR по
|
|
// доставке, с которой всё в порядке. Ошибка не логируется здесь — она
|
|
// возвращается наверх, и логирует её один раз вызывающий.
|
|
case <-ctx.Done():
|
|
log.Info("server stopping")
|
|
}
|
|
|
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), shutdownTimeout)
|
|
defer cancel()
|
|
|
|
// Приём прекращается РАНЬШЕ воркера: обратный порядок оставил бы доставки,
|
|
// принятые после его остановки, никого не разбудившими.
|
|
if err := srv.Shutdown(shutdownCtx); err != nil {
|
|
switch {
|
|
case errors.Is(err, context.DeadlineExceeded):
|
|
// Исчерпание бюджета Shutdown возвращает штатно, и отказом это не
|
|
// является: приём мог дочитывать многомегабайтное тело.
|
|
log.Warn("shutdown budget exceeded", "stage", "http")
|
|
case serveErr == nil:
|
|
serveErr = fmt.Errorf("shutdown: %w", err)
|
|
}
|
|
}
|
|
|
|
stopBackground()
|
|
if waitBackground(shutdownCtx, backgroundDone, log) {
|
|
closeStore()
|
|
}
|
|
return serveErr
|
|
}
|