Files
healthlog/internal/store/delivery.go
T
av 63bffe2865 Приём отвечает 200 до свёртки, свёртку ведёт фоновый воркер
- Очередью служит сама таблица: доставка ждёт свёртки в статусе `pending`,
  канал несёт только бит «есть работа». Переполнять нечего, падение процесса
  очередь не теряет, а подбор `pending` при старте — обычный проход воркера, а
  не отдельный код. Классификация исхода общая с пересборкой журнала.
- Исход разбора начал отражать доставку, а не обстоятельства: отмена и
  занятость базы статус не меняют (иначе конкуренция за базу выводила бы
  доставку из очереди навсегда), паника свёртки больше не валит процесс, а
  учёт доставки идёт через транзакцию с повторами.
- Длинный бюджет ответа выдан маршруту приёма, а не всему серверу:
  `write_timeout` в Go покрывает и чтение тела, и общий подъём снял бы защиту с
  остальных маршрутов.
2026-08-02 11:01:42 +03:00

364 lines
16 KiB
Go

package store
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"time"
)
// Статусы разбора доставки. Код ответа приёма от них не зависит: сохранили —
// значит приняли.
const (
// ParsePending — этим разбором тело ещё не смотрели.
//
// Смысл именно такой, а не «тела ещё не касались»: миграция 00005 перевела
// сюда доставки, разобранные кодом, который не различал частичный разбор.
// Статус консервативный — ретеншен не трогает pending никогда.
ParsePending = "pending"
// ParseDone — разобрано всё, что в теле было.
ParseDone = "parsed"
// ParsePartial — разобрано покрытое, но в теле остались секции, которых
// разбор не покрывает. Не отклонение, а установившееся состояние половины
// потока: 48 доставок из 99 несут только workouts или stateOfMind.
ParsePartial = "partial"
// ParseFailed — разобрать не удалось. Тело лежит в архиве, доставку
// подберёт пересборка.
ParseFailed = "failed"
)
// Delivery — учётная запись одного принятого пакета.
type Delivery struct {
ID string
ReceivedAt time.Time
AutomationName string
AutomationID string
Aggregation string
Period string
SessionID string
Bytes int64
SHA256 string
RawPath string
ParseStatus string
Points int64
// Headers — весь набор заголовков запроса как JSON-объект (имя → массив
// значений), уже без секретов. Именованные поля выше дублируют часть из
// них: по ним ходят запросы, а Headers хранит всё остальное на будущее.
Headers string
// UncoveredSections — секции тела, которых разбор не покрыл, JSON-массивом
// имён. Ответ на вопрос «что останется потерянным, если тело удалить»:
// ретеншен обязан спрашивать его прежде, чем срезать тело.
UncoveredSections string
}
// CreateDelivery записывает факт приёма пакета.
//
// Через ту же транзакцию с повторами, что и слияние точек, и это не симметрия
// ради симметрии. Свёртка держит запись всю доставку целиком — измерено 11
// секунд на 16 тысячах объектов, — а с фоновым воркером конкуренция за базу
// стала штатной. Одиночный `Exec` пересиживал бы только `busy_timeout`, после
// чего приём ответил бы `500` по доставке, тело которой уже на диске: доставка
// исчезла бы из журнала, а телефон её не перешлёт.
func (s *Store) CreateDelivery(ctx context.Context, d Delivery) error {
const q = `
INSERT INTO delivery (id, received_at, automation_name, automation_id,
aggregation, period, session_id, bytes, sha256,
raw_path, parse_status, points, headers)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`
headers := d.Headers
if headers == "" {
headers = "{}"
}
err := s.inTx(ctx, func(tx *sql.Tx) error {
_, err := tx.ExecContext(ctx, q,
d.ID, FormatTime(d.ReceivedAt), d.AutomationName, d.AutomationID,
d.Aggregation, d.Period, d.SessionID, d.Bytes, d.SHA256,
d.RawPath, d.ParseStatus, d.Points, headers)
return err //nolint:wrapcheck // обёртка одна, на выходе
})
if err != nil {
return fmt.Errorf("insert delivery: %w", err)
}
return nil
}
// LastDelivery возвращает последнюю по времени приёма доставку.
// Возвращает ErrNotFound, если доставок ещё не было.
func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) {
const q = `
SELECT id, received_at, automation_name, automation_id, aggregation,
period, session_id, bytes, sha256, raw_path, parse_status,
points, headers, uncovered_sections
FROM delivery ORDER BY received_at DESC, id DESC LIMIT 1`
var d Delivery
var receivedAt string
err := s.db.QueryRowxContext(ctx, q).Scan(
&d.ID, &receivedAt, &d.AutomationName, &d.AutomationID, &d.Aggregation,
&d.Period, &d.SessionID, &d.Bytes, &d.SHA256, &d.RawPath,
&d.ParseStatus, &d.Points, &d.Headers, &d.UncoveredSections)
if errors.Is(err, sql.ErrNoRows) {
return Delivery{}, ErrNotFound
}
if err != nil {
return Delivery{}, fmt.Errorf("select last delivery: %w", err)
}
d.ReceivedAt, err = ParseTime(receivedAt)
if err != nil {
return Delivery{}, err
}
return d, nil
}
// ListDeliveries возвращает учёт доставок в порядке журнала — `(received_at,
// id)`, тем же, в котором их проигрывает пересборка.
//
// Отдаются только **факты журнала**: то, что пришло вместе с доставкой.
// Производные от разбора поля (`parse_status`, `points`, `derived_layer`,
// `uncovered_sections`) сюда не попадают намеренно — перенос их в пересобранную
// базу сделал бы витрину функцией предыдущего прогона. Особенно `derived_layer`:
// доставка, чей повторный разбор отказал, отдала бы в наследование слой
// прежнего разбора, и следующая доставка той же автоматизации унаследовала бы
// его молча.
func (s *Store) ListDeliveries(ctx context.Context) ([]Delivery, error) {
const q = `
SELECT id, received_at, automation_name, automation_id, aggregation,
period, session_id, bytes, sha256, raw_path, headers
FROM delivery ORDER BY received_at, id`
rows, err := s.db.QueryContext(ctx, q)
if err != nil {
return nil, fmt.Errorf("select deliveries: %w", err)
}
defer func() { _ = rows.Close() }()
var out []Delivery
for rows.Next() {
var d Delivery
var receivedAt string
if err := rows.Scan(&d.ID, &receivedAt, &d.AutomationName, &d.AutomationID,
&d.Aggregation, &d.Period, &d.SessionID, &d.Bytes, &d.SHA256,
&d.RawPath, &d.Headers); err != nil {
return nil, fmt.Errorf("scan delivery: %w", err)
}
d.ReceivedAt, err = ParseTime(receivedAt)
if err != nil {
return nil, err
}
out = append(out, d)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("select deliveries: %w", err)
}
return out, nil
}
// PendingDelivery — доставка, ожидающая свёртки. Она же курсор обхода: место в
// журнале задаётся парой `(received_at, id)`, и вызывающему достаточно передать
// обратно последнюю полученную строку.
//
// Метка приёма отдаётся не для порядка (его держит SQL), а для метки отставания:
// «доставка ждала свёртки дольше N» считается от неё.
type PendingDelivery struct {
ID string
ReceivedAt time.Time
}
// PendingDeliveries возвращает неразобранные доставки в порядке журнала,
// строго после курсора. Нулевой курсор означает «с начала».
//
// Курсор нужен не ради страниц, а ради завершимости обхода: доставка, у которой
// не удалось записать даже исход разбора, остаётся `pending`, и выборка без
// курсора выдавала бы её бесконечно.
func (s *Store) PendingDeliveries(ctx context.Context, after PendingDelivery, limit int) ([]PendingDelivery, error) {
// Сравнение кортежем, а не через OR: развёрнутая форма даёт SCAN по
// индексу вместо SEARCH (проверено EXPLAIN QUERY PLAN). Тот же приём уже
// применён в LastDerivedLayer.
const q = `
SELECT id, received_at FROM delivery
WHERE parse_status = ? AND (received_at, id) > (?, ?)
ORDER BY received_at, id LIMIT ?`
rows, err := s.db.QueryContext(ctx, q, ParsePending, FormatTime(after.ReceivedAt), after.ID, limit)
if err != nil {
return nil, fmt.Errorf("select pending deliveries: %w", err)
}
defer func() { _ = rows.Close() }()
var out []PendingDelivery
for rows.Next() {
var d PendingDelivery
var receivedAt string
if err := rows.Scan(&d.ID, &receivedAt); err != nil {
return nil, fmt.Errorf("scan pending delivery: %w", err)
}
d.ReceivedAt, err = ParseTime(receivedAt)
if err != nil {
return nil, err
}
out = append(out, d)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("select pending deliveries: %w", err)
}
return out, nil
}
// CountPendingDeliveries возвращает размер задолженности — сколько доставок
// ждут свёртки. Нужен ровно одной строке лога при старте: сколько сервис должен
// разобрать, прежде чем витрина станет полной.
func (s *Store) CountPendingDeliveries(ctx context.Context) (int64, error) {
var n int64
if err := s.db.GetContext(ctx, &n,
`SELECT count(*) FROM delivery WHERE parse_status = ?`, ParsePending); err != nil {
return 0, fmt.Errorf("count pending deliveries: %w", err)
}
return n, nil
}
// DeliveryStatus возвращает статус разбора доставки.
func (s *Store) DeliveryStatus(ctx context.Context, id string) (string, error) {
var status string
err := s.db.GetContext(ctx, &status, `SELECT parse_status FROM delivery WHERE id = ?`, id)
if errors.Is(err, sql.ErrNoRows) {
return "", ErrNotFound
}
if err != nil {
return "", fmt.Errorf("select parse status: %w", err)
}
return status, nil
}
// FinishParse записывает исход разбора доставки.
//
// Слой сохраняется здесь же, потому что он нужен следующей доставке той же
// автоматизации: без плотных метрик выводить его не из чего, и наследовать
// приходится от прошлого раза.
//
// Пустой layer означает «не трогать»: у неудачной свёртки слоя нет, а
// затирание оборвало бы цепочку наследования — в том числе у доставки, которая
// раньше свернулась успешно, — и результат пересборки журнала изменился бы.
// ParseOutcome — исход разбора доставки. Структурой, а не растущим списком
// позиционных параметров: у FinishParse их было уже четыре, и пятый неизбежно
// перепутали бы местами с четвёртым.
type ParseOutcome struct {
Status string
Points int64
// Layer пустой означает «не трогать» — см. FinishParse.
Layer string
// Uncovered замещает прежнее значение ЦЕЛИКОМ, включая замещение пустым:
// у слоя пустота — отсутствие знания, у списка — знание об отсутствии.
// Пересвёртка доставки, чья секция стала покрытой, обязана список очистить.
Uncovered []string
}
func (s *Store) FinishParse(ctx context.Context, id string, out ParseOutcome) error {
const q = `
UPDATE delivery
SET parse_status = ?, points = ?,
derived_layer = CASE WHEN ? = '' THEN derived_layer ELSE ? END,
uncovered_sections = ?
WHERE id = ?`
// Ровно одно представление пустоты — `[]`: nil-срез Go сериализуется как
// null, и в колонке появилось бы второе значение с тем же смыслом.
sections := out.Uncovered
if sections == nil {
sections = []string{}
}
encoded, err := json.Marshal(sections)
if err != nil {
return fmt.Errorf("encode uncovered sections: %w", err)
}
res, err := s.db.ExecContext(ctx, q, out.Status, out.Points, out.Layer, out.Layer, string(encoded), id)
if err != nil {
return fmt.Errorf("update parse status: %w", err)
}
n, err := res.RowsAffected()
if err != nil {
return fmt.Errorf("update parse status: %w", err)
}
if n == 0 {
return ErrNotFound
}
return nil
}
// LastDerivedLayer возвращает слой, выведенный для этой автоматизации ПЕРЕД
// указанной доставкой. Пустая строка означает, что наследовать нечего.
//
// «Перед» здесь существенно, а не для красоты: слой обязан быть функцией от
// префикса журнала, иначе повторный прогон даёт другое состояние, чем живой
// приём. Запрос без границы по времени брал бы последний слой вообще — и при
// пересборке доставка наследовала бы слой от будущего. Проверено прогоном
// архива: 1737 объектов превращались в 1742.
func (s *Store) LastDerivedLayer(ctx context.Context, automationID string, before time.Time, beforeID string) (string, error) {
if automationID == "" {
return "", nil
}
const q = `
SELECT derived_layer FROM delivery
WHERE automation_id = ? AND derived_layer != ''
AND (received_at, id) < (?, ?)
ORDER BY received_at DESC, id DESC LIMIT 1`
var layer string
err := s.db.GetContext(ctx, &layer, q, automationID, FormatTime(before), beforeID)
if errors.Is(err, sql.ErrNoRows) {
return "", nil
}
if err != nil {
return "", fmt.Errorf("select derived layer: %w", err)
}
return layer, nil
}
// DeliveryBody — что нужно знать о доставке, чтобы разобрать её тело.
type DeliveryBody struct {
ID string
ReceivedAt time.Time
RawPath string
AutomationID string
Aggregation string
}
// DeliveryForParse возвращает сведения о доставке, нужные разбору.
func (s *Store) DeliveryForParse(ctx context.Context, id string) (DeliveryBody, error) {
const q = `
SELECT id, received_at, raw_path, automation_id, aggregation
FROM delivery WHERE id = ?`
var d DeliveryBody
var receivedAt string
err := s.db.QueryRowxContext(ctx, q, id).
Scan(&d.ID, &receivedAt, &d.RawPath, &d.AutomationID, &d.Aggregation)
if errors.Is(err, sql.ErrNoRows) {
return DeliveryBody{}, ErrNotFound
}
if err != nil {
return DeliveryBody{}, fmt.Errorf("select delivery for parse: %w", err)
}
d.ReceivedAt, err = ParseTime(receivedAt)
if err != nil {
return DeliveryBody{}, err
}
return d, nil
}
// CountDeliveries возвращает число принятых пакетов. Нужно для healthcheck и
// быстрой проверки «данные вообще идут».
func (s *Store) CountDeliveries(ctx context.Context) (int64, error) {
var n int64
if err := s.db.GetContext(ctx, &n, `SELECT count(*) FROM delivery`); err != nil {
return 0, fmt.Errorf("count deliveries: %w", err)
}
return n, nil
}