Files
healthlog/internal/ingest/ingest_test.go
T
av d33f37249c docs: документация переведена на канон av-dev-pm 3
- роадмап отвечает «что умеет и чего не умеет»: PLAN.md → ROADMAP.md, четыре
  канонические секции, достигнутые звенья строками в «Готово», цели
  переформулированы возможностями приложения
- задачи: род работы и «Затрагивает» набору спринта, 34 заголовка в форму
  действия, «Завершение» целей перечнями со ссылкой из каждой задачи
- вычитка проходами task-form и doc-wording, починены протухшие факты в README,
  паспорте и review.md
2026-08-04 20:48:30 +03:00

334 lines
12 KiB
Go

package ingest_test
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"io"
"log/slog"
"os"
"path/filepath"
"strings"
"testing"
"git.vakhrushev.me/av/healthlog/internal/archive"
"git.vakhrushev.me/av/healthlog/internal/ingest"
"git.vakhrushev.me/av/healthlog/internal/store"
)
func TestAcceptStoresBodyVerbatim(t *testing.T) {
svc, arch, st := newService(t)
// Незнакомая метрика и неизвестная форма точки — не повод отвергать пакет.
body := []byte(`{"data":{"metrics":[{"name":"новая_метрика","units":"?","data":[{"чего_то":1,"date":"2026-07-31 12:00:00 +0300"}]}]}}`)
res, err := svc.Accept(context.Background(), body, ingest.Meta{
AutomationName: "всё подряд",
Aggregation: "None",
Period: "Since Last Sync",
})
if err != nil {
t.Fatalf("Accept: %v", err)
}
if res.Bytes != int64(len(body)) {
t.Errorf("bytes = %d, ожидалось %d", res.Bytes, len(body))
}
sum := sha256.Sum256(body)
if res.SHA256 != hex.EncodeToString(sum[:]) {
t.Errorf("sha256 = %q", res.SHA256)
}
// Тело в архиве — байт в байт то, что прислали.
r, err := arch.Open(res.RawPath)
if err != nil {
t.Fatalf("Open: %v", err)
}
defer func() { _ = r.Close() }()
got, err := io.ReadAll(r)
if err != nil {
t.Fatalf("чтение архива: %v", err)
}
if string(got) != string(body) {
t.Errorf("в архиве %q, прислано %q", got, body)
}
n, err := st.CountDeliveries(context.Background())
if err != nil {
t.Fatalf("CountDeliveries: %v", err)
}
if n != 1 {
t.Errorf("доставок в БД: %d, ожидалась 1", n)
}
}
// Повтор того же тела пока принимается — отсев идентичных доставок отложен
// (docs/tasks/ROADMAP.md). Проверяем, что повтор не ломается и не затирает первую.
func TestAcceptAllowsRepeatedBody(t *testing.T) {
svc, _, st := newService(t)
body := []byte(`{"data":{"metrics":[]}}`)
first, err := svc.Accept(context.Background(), body, ingest.Meta{})
if err != nil {
t.Fatalf("первый Accept: %v", err)
}
second, err := svc.Accept(context.Background(), body, ingest.Meta{})
if err != nil {
t.Fatalf("второй Accept: %v", err)
}
if first.DeliveryID == second.DeliveryID {
t.Error("у двух доставок совпали идентификаторы")
}
if first.SHA256 != second.SHA256 {
t.Error("у одинаковых тел разошёлся sha256")
}
n, err := st.CountDeliveries(context.Background())
if err != nil {
t.Fatalf("CountDeliveries: %v", err)
}
if n != 2 {
t.Errorf("доставок в БД: %d, ожидалось 2", n)
}
}
func TestAcceptRejectsMalformed(t *testing.T) {
cases := map[string]string{
"пустое тело": ``,
"не JSON": `не json`,
"обрезанный JSON": `{"data":{"metrics":[`,
"нет data": `{"metrics":[]}`,
"data не объект": `{"data":[]}`,
}
for name, body := range cases {
t.Run(name, func(t *testing.T) {
svc, _, st := newService(t)
_, err := svc.Accept(context.Background(), []byte(body), ingest.Meta{})
if !errors.Is(err, ingest.ErrMalformed) {
t.Fatalf("ошибка = %v, ожидалась ErrMalformed", err)
}
// Отвергнутая доставка не оставляет следа в учёте.
n, err := st.CountDeliveries(context.Background())
if err != nil {
t.Fatalf("CountDeliveries: %v", err)
}
if n != 0 {
t.Errorf("доставок в БД: %d, ожидалось 0", n)
}
})
}
}
// Ответ отдаётся ДО свёртки: принятая доставка ждёт разбора в очереди, а не
// приезжает разобранной. Это смена контракта, и она проверяется явно.
func TestAcceptОставляетДоставкуВОчереди(t *testing.T) {
svc, _, st := newService(t)
ctx := context.Background()
body, err := os.ReadFile(filepath.Join("..", "hae", "testdata", "minute.json"))
if err != nil {
t.Fatalf("фикстура: %v", err)
}
if _, err := svc.Accept(ctx, body, ingest.Meta{Aggregation: "Minutes"}); err != nil {
t.Fatalf("Accept: %v", err)
}
d, err := st.LastDelivery(ctx)
if err != nil {
t.Fatalf("LastDelivery: %v", err)
}
if d.ParseStatus != store.ParsePending {
t.Errorf("parse_status = %q, ожидался %q", d.ParseStatus, store.ParsePending)
}
n, err := st.CountBuckets(ctx)
if err != nil {
t.Fatalf("CountBuckets: %v", err)
}
if n != 0 {
t.Errorf("объектов %d: свёртка произошла внутри приёма", n)
}
}
// Сигнал уходит после того, как доставка учтена: воркер, разбуженный раньше,
// потратил бы проход впустую.
func TestAcceptБудитСвёрткуПослеУчёта(t *testing.T) {
dir := t.TempDir()
st, arch := newDeps(t, dir)
var seen int64
notify := func() {
n, err := st.CountDeliveries(context.Background())
if err != nil {
t.Errorf("CountDeliveries: %v", err)
}
seen = n
}
svc := ingest.New(arch, st, notify, slog.New(slog.DiscardHandler))
if _, err := svc.Accept(context.Background(), []byte(`{"data":{"metrics":[]}}`), ingest.Meta{}); err != nil {
t.Fatalf("Accept: %v", err)
}
if seen != 1 {
t.Errorf("на момент сигнала доставок в учёте %d, ожидалась 1", seen)
}
}
// Нулевой сигнал — законный вход (свёрткой заведует вызывающий), и приём от
// него не падает. Паника здесь наступила бы ПОСЛЕ записи тела и учёта, то есть
// отправитель получил бы отказ по сохранённой доставке.
func TestAcceptБезСигналаНеПадает(t *testing.T) {
dir := t.TempDir()
st, arch := newDeps(t, dir)
svc := ingest.New(arch, st, nil, slog.New(slog.DiscardHandler))
if _, err := svc.Accept(context.Background(), []byte(`{"data":{"metrics":[]}}`), ingest.Meta{}); err != nil {
t.Fatalf("Accept: %v", err)
}
}
// Обрыв соединения после записи тела не должен оставлять тело без учёта:
// доставка, не попавшая в журнал, восстанавливается только пересборкой с
// ручной подменой базы.
func TestAcceptУчитываетДоставкуПослеОбрываСоединения(t *testing.T) {
svc, _, st := newService(t)
ctx, cancel := context.WithCancel(context.Background())
cancel()
res, err := svc.Accept(ctx, []byte(`{"data":{"metrics":[]}}`), ingest.Meta{})
if err != nil {
t.Fatalf("Accept на отменённом контексте: %v", err)
}
status, err := st.DeliveryStatus(context.Background(), res.DeliveryID)
if err != nil {
t.Fatalf("DeliveryStatus: %v", err)
}
if status != store.ParsePending {
t.Errorf("parse_status = %q, ожидался %q", status, store.ParsePending)
}
}
// Учёта нет, а тело есть: приём кладёт тело на диск раньше строки в базе, и
// отказ на вставке оставляет тело в архиве. Такое тело подберёт пересборка.
func TestAcceptПриОтказеУчётаОставляетТелоВАрхиве(t *testing.T) {
dir := t.TempDir()
st, arch := newDeps(t, dir)
svc := ingest.New(arch, st, nil, slog.New(slog.DiscardHandler))
// База закрыта — учесть доставку нечем.
if err := st.Close(); err != nil {
t.Fatalf("закрытие базы: %v", err)
}
if _, err := svc.Accept(context.Background(), []byte(`{"data":{"metrics":[]}}`), ingest.Meta{}); err == nil {
t.Fatal("приём не заметил, что доставка не учтена")
}
entries, err := filepath.Glob(filepath.Join(arch.Root(), "*", "*", "*", "*.json.gz"))
if err != nil {
t.Fatalf("обход архива: %v", err)
}
if len(entries) != 1 {
t.Errorf("тел в архиве %d, ожидалось 1: тело потеряно вместе с учётом", len(entries))
}
}
func newDeps(t *testing.T, dir string) (*store.Store, *archive.Archive) {
t.Helper()
st, err := store.Open(filepath.Join(dir, "healthlog.db"))
if err != nil {
t.Fatalf("store.Open: %v", err)
}
t.Cleanup(func() { _ = st.Close() })
arch, err := archive.New(filepath.Join(dir, "raw"))
if err != nil {
t.Fatalf("archive.New: %v", err)
}
return st, arch
}
func newService(t *testing.T) (*ingest.Service, *archive.Archive, *store.Store) {
t.Helper()
st, arch := newDeps(t, t.TempDir())
return ingest.New(arch, st, nil, slog.New(slog.DiscardHandler)), arch, st
}
// Отказ учёта после того, как тело легло в архив, обязан называть класс
// причины: занятость базы — конкуренция за запись, которая будет повторяться, и
// лечится она не тем же, чем сбой диска. Уровень при этом остаётся ERROR: тело
// осиротело в любом случае, и вернуть его в журнал может только пересборка.
func TestAcceptОтказУчётаНазываетКлассПричины(t *testing.T) {
dir := t.TempDir()
st, arch := newDeps(t, dir)
var buf bytes.Buffer
log := slog.New(slog.NewJSONHandler(&buf, nil))
svc := ingest.New(arch, st, nil, log)
if err := st.Close(); err != nil {
t.Fatalf("закрытие базы: %v", err)
}
if _, err := svc.Accept(context.Background(), []byte(`{"data":{"metrics":[]}}`), ingest.Meta{}); err == nil {
t.Fatal("приём не заметил, что доставка не учтена")
}
var rec map[string]any
for _, line := range strings.Split(strings.TrimSpace(buf.String()), "\n") {
var v map[string]any
if err := json.Unmarshal([]byte(line), &v); err != nil {
t.Fatalf("строка лога не JSON: %v", err)
}
if v["msg"] == "delivery failed" {
rec = v
}
}
if rec == nil {
t.Fatal("отказ учёта не залогирован")
}
if rec["level"] != "ERROR" {
t.Errorf("уровень %v, ожидался ERROR: тело осиротело", rec["level"])
}
busy, ok := rec["db_busy"].(bool)
if !ok {
t.Fatalf("класс причины не назван: %v", rec)
}
// Закрытая база — не занятость: признак обязан различать, а не стоять всегда.
if busy {
t.Error("закрытая база названа занятой — признак не различает причины")
}
}
// Инвариант «тела запросов только на DEBUG и с обрезкой» относится и к DEBUG:
// проверка формы конверта идёт через encoding/json, чей UnmarshalTypeError
// кладёт в текст литерал значения.
func TestAcceptОтказФормыНеНесётТелаВЛог(t *testing.T) {
var buf bytes.Buffer
log := slog.New(slog.NewJSONHandler(&buf, &slog.HandlerOptions{Level: slog.LevelDebug}))
// Ни архив, ни база не нужны: тело неверной формы отвергается проверкой
// конверта до всякой записи.
svc := ingest.New(nil, nil, nil, log)
body := []byte(`{"data":` + strings.Repeat("9", 1<<20) + `}`)
if _, err := svc.Accept(context.Background(), body, ingest.Meta{}); err == nil {
t.Fatal("тело неверной формы принято")
}
if buf.Len() > 4096 {
t.Errorf("строка лога %d Б: содержимое тела уехало в лог", buf.Len())
}
if strings.Contains(buf.String(), strings.Repeat("9", 256)) {
t.Error("литерал из тела виден в логе")
}
}