- internal/fold — свёртка по идентификатору доставки, тело из архива: тот же код, каким пойдёт пересборка витрины - приём зовёт свёртку на context.WithoutCancel с собственным дедлайном; исход разбора на код ответа не влияет - слой доставки хранится в delivery.derived_layer и наследуется от ПРЕДШЕСТВУЮЩЕЙ доставки автоматизации: без границы по времени свёртка переставала быть функцией от префикса журнала (1737 объектов против 1742) - task verify:archive — сходимость на живом архиве, 99 доставок из 99
185 lines
8.0 KiB
Go
185 lines
8.0 KiB
Go
// Package ingest — use-case приёма пакета от Health Auto Export, общий для
|
||
// HTTP и (позже) CLI-импорта.
|
||
package ingest
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
"crypto/sha256"
|
||
"encoding/hex"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"log/slog"
|
||
"time"
|
||
|
||
"git.vakhrushev.me/av/healthlog/internal/archive"
|
||
"git.vakhrushev.me/av/healthlog/internal/fold"
|
||
"git.vakhrushev.me/av/healthlog/internal/ident"
|
||
"git.vakhrushev.me/av/healthlog/internal/store"
|
||
)
|
||
|
||
// ErrMalformed — тело не разбирается как JSON ожидаемой формы. Это отказ
|
||
// доставки (обрыв, обрезанное тело), о нём отправителю сообщаем ошибкой.
|
||
// Незнакомое СОДЕРЖИМОЕ (новая метрика, новая форма точки) сюда не относится:
|
||
// такой пакет принимается, сохраняется и доразбирается позже.
|
||
var ErrMalformed = errors.New("некорректный формат пакета")
|
||
|
||
// Meta — что рассказала о себе автоматизация Health Auto Export (свои
|
||
// заголовки запроса).
|
||
type Meta struct {
|
||
AutomationName string
|
||
AutomationID string
|
||
Aggregation string
|
||
Period string
|
||
SessionID string
|
||
// Headers — весь набор заголовков запроса, уже без секретов (вырезает
|
||
// транспорт: только он знает список токенов). Документация HAE называет
|
||
// пять заголовков, но что приложение шлёт на самом деле — неизвестно, и
|
||
// именно из незадокументированного вышли самые полезные находки.
|
||
Headers map[string][]string
|
||
}
|
||
|
||
// Result — итог приёма.
|
||
type Result struct {
|
||
DeliveryID string
|
||
Bytes int64
|
||
SHA256 string
|
||
RawPath string
|
||
}
|
||
|
||
// foldTimeout — сколько отводится свёртке принятой доставки.
|
||
//
|
||
// Свёртка идёт на контексте, отвязанном от запроса, поэтому собственный
|
||
// дедлайн обязателен: без него зависшая запись держала бы горутину до конца
|
||
// жизни процесса.
|
||
const foldTimeout = 2 * time.Minute
|
||
|
||
// Service принимает пакеты: сохраняет тело в архив, учитывает доставку и
|
||
// запускает её свёртку.
|
||
type Service struct {
|
||
arch *archive.Archive
|
||
store *store.Store
|
||
fold *fold.Service
|
||
log *slog.Logger
|
||
}
|
||
|
||
// New собирает use-case приёма.
|
||
func New(arch *archive.Archive, st *store.Store, f *fold.Service, log *slog.Logger) *Service {
|
||
return &Service{arch: arch, store: st, fold: f, log: log.With("capability", "ingest")}
|
||
}
|
||
|
||
// Accept принимает тело пакета: проверяет форму, кладёт в сырой архив и
|
||
// заводит запись о доставке.
|
||
//
|
||
// Порядок важен: сначала тело оказывается на диске, и только потом появляется
|
||
// учётная запись. Обратный порядок дал бы учтённую доставку без данных.
|
||
//
|
||
// Это единственный логирующий чекпоинт приёма — транспорт исход не логирует.
|
||
func (s *Service) Accept(ctx context.Context, body []byte, meta Meta) (Result, error) {
|
||
if err := checkEnvelope(body); err != nil {
|
||
// Некорректный ввод от отправителя — норма жизни, команде разбирать
|
||
// нечего: DEBUG, а не ERROR.
|
||
s.log.DebugContext(ctx, "delivery rejected", "error", err, "bytes", len(body))
|
||
return Result{}, err
|
||
}
|
||
|
||
sum := sha256.Sum256(body)
|
||
res := Result{
|
||
DeliveryID: ident.NewID(),
|
||
Bytes: int64(len(body)),
|
||
SHA256: hex.EncodeToString(sum[:]),
|
||
}
|
||
receivedAt := store.Now()
|
||
|
||
rawPath, err := s.arch.Write(res.DeliveryID, receivedAt, body)
|
||
if err != nil {
|
||
s.log.ErrorContext(ctx, "delivery failed", "error", err, "delivery_id", res.DeliveryID)
|
||
return Result{}, fmt.Errorf("archive body: %w", err)
|
||
}
|
||
res.RawPath = rawPath
|
||
|
||
err = s.store.CreateDelivery(ctx, store.Delivery{
|
||
ID: res.DeliveryID,
|
||
ReceivedAt: receivedAt,
|
||
Headers: encodeHeaders(meta.Headers),
|
||
AutomationName: meta.AutomationName,
|
||
AutomationID: meta.AutomationID,
|
||
Aggregation: meta.Aggregation,
|
||
Period: meta.Period,
|
||
SessionID: meta.SessionID,
|
||
Bytes: res.Bytes,
|
||
SHA256: res.SHA256,
|
||
RawPath: rawPath,
|
||
ParseStatus: store.ParsePending,
|
||
})
|
||
if err != nil {
|
||
// Тело уже на диске — данные не потеряны, но учёта нет. Разбор архива
|
||
// на следующем шаге проекта такую доставку подберёт.
|
||
s.log.ErrorContext(ctx, "delivery failed", "error", err, "delivery_id", res.DeliveryID, "raw_path", rawPath)
|
||
return Result{}, fmt.Errorf("record delivery: %w", err)
|
||
}
|
||
|
||
s.log.InfoContext(ctx, "delivery accepted",
|
||
"delivery_id", res.DeliveryID,
|
||
"bytes", res.Bytes,
|
||
"raw_path", rawPath,
|
||
"automation_name", meta.AutomationName,
|
||
"aggregation", meta.Aggregation,
|
||
"period", meta.Period)
|
||
|
||
// Свёртка идёт после того, как доставка учтена, и на контексте, ОТВЯЗАННОМ
|
||
// от запроса: обрыв соединения клиентом или прокси на середине оставил бы
|
||
// часть объектов записанной, а доставку — со статусом, по которому её
|
||
// никто не подберёт. Исход свёртки на код ответа не влияет — сохранили
|
||
// значит приняли.
|
||
foldCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), foldTimeout)
|
||
defer cancel()
|
||
// Ошибку не возвращаем: она уже записана в лог и в parse_status свёрткой,
|
||
// а доставка принята.
|
||
_, _ = s.fold.Fold(foldCtx, res.DeliveryID)
|
||
|
||
return res, nil
|
||
}
|
||
|
||
// encodeHeaders сериализует заголовки для хранения. Ключи json.Marshal
|
||
// сортирует сам, поэтому запись стабильна и её удобно сравнивать между
|
||
// доставками. Сбой сериализации не должен ронять приём: заголовки —
|
||
// вспомогательные сведения, а не данные, ради которых всё затевалось.
|
||
func encodeHeaders(h map[string][]string) string {
|
||
if len(h) == 0 {
|
||
return "{}"
|
||
}
|
||
b, err := json.Marshal(h)
|
||
if err != nil {
|
||
return "{}"
|
||
}
|
||
return string(b)
|
||
}
|
||
|
||
// envelope — минимальная форма пакета Health Auto Export.
|
||
type envelope struct {
|
||
Data json.RawMessage `json:"data"`
|
||
}
|
||
|
||
// checkEnvelope проверяет ровно то, что делает пакет доставленным: это JSON,
|
||
// в нём есть объект data. Содержимое data не трогаем — его форму определяет
|
||
// Apple, и незнакомая форма не повод отвергать доставку.
|
||
func checkEnvelope(body []byte) error {
|
||
if len(bytes.TrimSpace(body)) == 0 {
|
||
return fmt.Errorf("%w: пустое тело", ErrMalformed)
|
||
}
|
||
|
||
var env envelope
|
||
if err := json.Unmarshal(body, &env); err != nil {
|
||
return fmt.Errorf("%w: %v", ErrMalformed, err) //nolint:errorlint // причину наружу не раскрываем, она уходит в лог
|
||
}
|
||
if len(env.Data) == 0 {
|
||
return fmt.Errorf("%w: нет объекта data", ErrMalformed)
|
||
}
|
||
if trimmed := bytes.TrimSpace(env.Data); len(trimmed) == 0 || trimmed[0] != '{' {
|
||
return fmt.Errorf("%w: data не объект", ErrMalformed)
|
||
}
|
||
return nil
|
||
}
|