- половина потока (50 доставок из 104) не несёт metrics вовсе и до сих пор числилась parsed: ретеншен, поверив статусу, срезал бы тела stateOfMind, которых в экспорте Apple нет - разбор перечисляет верхнеуровневые ключи data, непокрытые проглатываются декодированием: тело 40 МиБ из непокрытой секции удерживает 0 МиБ - статус partial и колонка delivery.uncovered_sections; миграция переводит прежние parsed в pending — им верить нельзя - витрина не изменилась: отпечаток совпал с прогоном до изменения
235 lines
9.9 KiB
Go
235 lines
9.9 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 записывает факт приёма пакета.
|
|
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.db.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)
|
|
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
|
|
}
|
|
|
|
// 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
|
|
}
|