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 // SkippedEntities — сколько сущностей с собственным `id` разбор пропустил. // Второй половина ответа на тот же вопрос: сущность, которую разбор не // понял, в витрину не попала, а список непокрытых секций про неё молчит. // // Отсутствие значения означает «не измерялось» и НЕ равно нулю: так // выглядят доставки, свёрнутые разбором, который пропусков не считал, и те, // чей разбор не досчитал. Читатель, принимающий по счётчику необратимое // решение, обязан трактовать отсутствие как «не удалять». // // Указателем, а не sql.NullInt64: поле уедет в JSON `/stats` и в MCP, а // NullInt64 сериализуется формой драйвера (`{"Int64":0,"Valid":false}`) — // первый, кто про это забудет, опубликует её наружу, и она станет // контрактом. Указатель даёт `null` бесплатно и означает ровно то же. SkippedEntities *int64 } // 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, skipped_entities FROM delivery ORDER BY received_at DESC, id DESC LIMIT 1` var d Delivery var receivedAt string // sql.NullInt64 живёт ровно на границе сканирования и наружу не выходит. var skipped sql.NullInt64 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, &skipped) if errors.Is(err, sql.ErrNoRows) { return Delivery{}, ErrNotFound } if err != nil { return Delivery{}, fmt.Errorf("select last delivery: %w", err) } if skipped.Valid { d.SkippedEntities = &skipped.Int64 } 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`, `skipped_entities`) сюда не попадают намеренно — // перенос их в пересобранную базу сделал бы витрину функцией предыдущего // прогона. Особенно `derived_layer`: доставка, чей повторный разбор отказал, // отдала бы в наследование слой прежнего разбора, и следующая доставка той же // автоматизации унаследовала бы его молча. У `skipped_entities` цена та же и // хуже: пустота у него значит «не измерялось», и перенесённое число выдавало бы // измерение прежнего разбора за измерение текущего — а по нему принимается // необратимое решение об удалении тела. // // Перечень пополняется ТЕМ ЖЕ изменением, которое заводит новое поле: он // единственное место, где сказано, чему нельзя пережить пересборку, и следующий // автор решает по нему. Поле, не внесённое сюда, однажды перенесут «для полноты // учёта». 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 } // ParsedAfter отвечает, есть ли в журнале доставка ПОЗЖЕ указанной позиции, // уже записавшая свой разбор в витрину. // // Нужен ровно одному наблюдению: свёртка вне порядка журнала. Правило слияния // точек разрешает равную полноту в пользу пришедшей доставки, то есть исход // есть функция порядка свёртки; доставка, свёрнутая после своей преемницы, // возвращает координату к версии, которую источник уже пересчитал, и живая // витрина расходится с пересборкой. Барьер воркера этого не ловит: отставшей // доставки в момент прохода просто не существует — строка учёта становится // видимой только после записи тела. // // `failed` в предикат НЕ входит: такая доставка в витрину ничего не записала, // перестановка относительно неё содержимого не разводит, а исход этот штатный // (невыводимый слой), и его учёт превратил бы WARN в шум. // // Это страж окна, а не постоянная часть свёртки: закрыв порядок на самом // приёме, метод и его вызов сносят вместе с окном. func (s *Store) ParsedAfter(ctx context.Context, at time.Time, id string) (bool, error) { // Кортежное сравнение, а не через OR: развёрнутая форма даёт SCAN по // индексу вместо SEARCH — тот же приём, что в PendingDeliveries. const q = ` SELECT EXISTS( SELECT 1 FROM delivery WHERE parse_status IN (?, ?) AND (received_at, id) > (?, ?))` var found bool if err := s.db.GetContext(ctx, &found, q, ParseDone, ParsePartial, FormatTime(at), id); err != nil { return false, fmt.Errorf("select parsed after: %w", err) } return found, 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 // SkippedEntities — сколько сущностей с собственным `id` разбор пропустил. // Пишется, когда разбор ДОСЧИТАЛ, включая ноль: доставка, пропуски которой // исчезли вместе с поумневшим разбором, не должна остаться помеченной // навсегда. // // nil означает «не измерялось» и колонку НЕ ТРОГАЕТ — та же идиома, что у // пустого Layer. Без неё отказ на чтении тела или паника разбора писали бы // ноль, то есть «проверено, терять нечего», в доставку, содержимое которой // никто не смотрел: ровно та подстановка, ради отказа от которой колонка // заведена без DEFAULT. SkippedEntities *int64 } 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 = ?, skipped_entities = CASE WHEN ? THEN ? ELSE skipped_entities END 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) } // Отсутствие числа не пишется нулём: ноль означает «измерено, пропусков не // было», а нам нужно «не измерялось». Колонка остаётся какой была — та же // форма, что у слоя строкой выше. measured := out.SkippedEntities != nil var skipped int64 if measured { skipped = *out.SkippedEntities } res, err := s.db.ExecContext(ctx, q, out.Status, out.Points, out.Layer, out.Layer, string(encoded), measured, skipped, 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 // Headers — заголовки запроса JSON-объектом. Разбору нужен ровно один из // них, `Accept-Language`: он задаёт язык, на котором приехали // категориальные строки, и сужает поиск по словарю. Колонки под язык нет // намеренно — она была бы вторым домом факта, который уже лежит здесь. Headers string } // DeliveryForParse возвращает сведения о доставке, нужные разбору. func (s *Store) DeliveryForParse(ctx context.Context, id string) (DeliveryBody, error) { const q = ` SELECT id, received_at, raw_path, automation_id, aggregation, headers 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, &d.Headers) 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 }