httpapi: точки метрики за период отдаются одним запросом

- `GET /api/v1/metrics/{name}?from&to&layer` — ряд точек за период; конверт
  объявляет слой, измеренный род, его применимость к отданному ряду и границу
  окна измерения, а сам ряд собирается из одного слоя, выбранного по охвату
  точек внутри периода
- use-case вынесен в `internal/points`, чтение — одним входом `store.ReadSeries`
  под одной транзакцией; правило выбора слоя остаётся в домене и приходит в
  хранилище колбэком
- `writeJSON` перестал экранировать HTML-символы и перестал глушить отказ
  записи: дословность содержимого точки иначе не удерживается, а оборванное
  тело уходило под видом успешного `200`
This commit is contained in:
av
2026-08-04 18:46:45 +03:00
parent b819b77f62
commit 29ca8d415c
36 changed files with 4721 additions and 58 deletions
+40 -16
View File
@@ -147,7 +147,12 @@ func closeEnough(a, b float64) bool {
return math.Abs(a-b)/math.Max(math.Abs(a), math.Abs(b)) <= tolerance
}
func clipMetric(metric string) string {
// ClipMetric обрезает имя метрики для записи лога.
//
// Экспортирована потому, что предел один на всех, кто пишет имя метрики в лог:
// имя приходит из тела дословно при пределе приёма в 64 МиБ, а запись
// повторяется на каждый запрос. Второй предел разошёлся бы с первым молча.
func ClipMetric(metric string) string {
if len(metric) <= maxMetricInLog {
return metric
}
@@ -252,12 +257,38 @@ func (s *Service) Version(ctx context.Context) (string, error) {
s.log.DebugContext(ctx, "state version unavailable", "capability", "query", "error", err)
return "", err
}
return stamp(version, store.Now().Add(horizonSlack)), nil
return Stamp(version, Horizon()), nil
}
// stamp склеивает версию витрины с горизонтом. Пустая версия остаётся пустой:
// Horizon — верхняя граница окна измерения на текущий момент.
//
// Экспортирована потому, что горизонт нужен ВСЕМ, кто объявляет измеренный род:
// маршрут точек снимает род и метку с одного горизонта, иначе метка подтвердит
// неизменность ответа, в котором род уже перевернулся ходом часов.
func Horizon() time.Time { return store.Now().Add(horizonSlack) }
// MeasureWindow — окно измерения рода для заданного горизонта.
//
// Одно на всех потребителей измерения. Второй экземпляр параметров разошёлся бы
// с первым молча, а вердикт зависит от каждого из них: слои сверки, размер окна
// и оба структурных порога уходят в предварительный отбор хранилища.
func MeasureWindow(horizon time.Time) store.CatalogWindow {
return store.CatalogWindow{
Fine: string(hae.LayerMinute),
Coarse: string(hae.LayerHour),
Hours: Window,
Horizon: horizon,
CoarsePoints: coarsePoints,
MinFinePoints: minFinePoints,
}
}
// Stamp склеивает версию витрины с горизонтом. Пустая версия остаётся пустой:
// подписывать нечем — значит нечем, и горизонт этого не меняет.
func stamp(version string, horizon time.Time) string {
//
// Экспортирована по той же причине, что и Horizon: правило «метка строится из
// всего, от чего зависит ответ» держится ровно до тех пор, пока склейка одна.
func Stamp(version string, horizon time.Time) string {
if version == "" {
return ""
}
@@ -287,19 +318,12 @@ func New(st *store.Store, log *slog.Logger) *Service {
// её хранилище — двумя пробами вокруг чтения. Порядок проб там же и объяснён:
// версия, снятая после чтения, пометила бы устаревший снимок свежей меткой.
func (s *Service) Metrics(ctx context.Context) (Snapshot, error) {
horizon := store.Now().Add(horizonSlack)
horizon := Horizon()
var snap store.CatalogSnapshot
version, err := s.store.VersionedRead(ctx, func(ctx context.Context) error {
var err error
snap, err = s.store.ReadCatalog(ctx, store.CatalogWindow{
Fine: string(hae.LayerMinute),
Coarse: string(hae.LayerHour),
Hours: Window,
Horizon: horizon,
CoarsePoints: coarsePoints,
MinFinePoints: minFinePoints,
})
snap, err = s.store.ReadCatalog(ctx, MeasureWindow(horizon))
return err
})
if err != nil { //nolint:nestif // ветка одна, вложенность даёт лог по адресату
@@ -331,7 +355,7 @@ func (s *Service) Metrics(ctx context.Context) (Snapshot, error) {
if basis.Conflicting > 0 {
s.log.WarnContext(ctx, "aggregation style conflict",
"capability", "query",
"metric", clipMetric(group.metric),
"metric", ClipMetric(group.metric),
"hours", basis.Hours,
"compared", basis.Compared,
"agreeing", basis.Agreeing,
@@ -344,7 +368,7 @@ func (s *Service) Metrics(ctx context.Context) (Snapshot, error) {
if to := group.latest(); to.After(horizon) {
s.log.WarnContext(ctx, "future data",
"capability", "query",
"metric", clipMetric(group.metric),
"metric", ClipMetric(group.metric),
"last_ts", store.FormatTime(to),
"horizon", store.FormatTime(horizon))
}
@@ -362,7 +386,7 @@ func (s *Service) Metrics(ctx context.Context) (Snapshot, error) {
// механизм не окупается вовсе.
s.log.DebugContext(ctx, "catalog unsigned", "capability", "query")
}
return Snapshot{Version: stamp(version, horizon), Metrics: out}, nil
return Snapshot{Version: Stamp(version, horizon), Metrics: out}, nil
}
type metricGroup struct {
+26
View File
@@ -520,3 +520,29 @@ func TestКаталогОтдаётсяСВерсиейВитрины(t *testing
t.Error("каталог собран на стоящей витрине и остался без версии")
}
}
// Версия ответа каталога — это версия витрины ПЛЮС горизонт измерения.
//
// Утверждение прямое, потому что склейка теперь общая: её же зовёт маршрут
// точек. Сломай её — и оба маршрута начнут подтверждать неизменность ответа,
// чей род перевернулся ходом часов, а не коммитом.
func TestВерсияКаталогаНесётГоризонт(t *testing.T) {
st := openStore(t)
ctx := context.Background()
got, err := service(t, st).Version(ctx)
if err != nil {
t.Fatalf("Version: %v", err)
}
bare, err := st.StateVersion(ctx)
if err != nil {
t.Fatalf("StateVersion: %v", err)
}
if got == bare {
t.Error("версия ответа равна версии витрины — горизонт в неё не вошёл")
}
if want := catalog.Stamp(bare, catalog.Horizon()); got != want {
t.Errorf("версия ответа %q, ожидалась %q", got, want)
}
}
+3 -3
View File
@@ -15,16 +15,16 @@ func TestГоризонтВходитВВерсиюОтвета(t *testing.T) {
at := time.Date(2026, 6, 1, 10, 30, 0, 0, time.UTC)
if stamp("v", at) == stamp("v", at.Add(2*time.Hour)) {
if Stamp("v", at) == Stamp("v", at.Add(2*time.Hour)) {
t.Error("версия не изменилась при сдвиге горизонта на два часа")
}
// Огрубление до часа точное, а не приблизительное: `hour_utc` объектов лежит
// ровно на часах, поэтому отбор меняется ровно при переходе через час.
// Внутри часа метка обязана стоять — иначе она дребезжала бы ежесекундно.
if stamp("v", at) != stamp("v", at.Add(20*time.Minute)) {
if Stamp("v", at) != Stamp("v", at.Add(20*time.Minute)) {
t.Error("версия сдвинулась внутри одного часа — метка дребезжит на месте")
}
if stamp("", at) != "" {
if Stamp("", at) != "" {
t.Error("пустая версия витрины подписана горизонтом — подписывать нечем")
}
}
+26 -15
View File
@@ -856,25 +856,36 @@ func headerLayer(aggregation string) Layer {
// finer возвращает более мелкий из двух слоёв.
func finer(a, b Layer) Layer {
if rank(a) < rank(b) {
if Rank(a) < Rank(b) {
return a
}
return b
}
func rank(l Layer) int {
switch l {
case LayerSample:
return 0
case LayerRaw:
return 1
case LayerMinute:
return 2
case LayerHour:
return 3
case LayerDay:
return 4
default:
return 5
// Layers — ЕДИНСТВЕННЫЙ словарь слоёв, упорядоченный от самого мелкого к самому
// крупному.
//
// Один, потому что иначе их становится четыре: порядок для разбора, перечень
// для выборки охватов, проверка параметра запроса и текст отказа клиенту. Ни
// компилятор, ни тест их не сверяют — новая константа `Layer` скомпилировалась
// бы, получила бы ранг «крупнее всех», не попала бы в выборку охватов и
// отвергалась бы маршрутом как незнакомая. Симптомом был бы пустой ряд при
// непустых данных. Случай не гипотетический: слой `sample` наполнится импортом
// родного экспорта Apple.
//
// Срез, а не карта: порядок здесь и есть содержание.
var Layers = []Layer{LayerSample, LayerRaw, LayerMinute, LayerHour, LayerDay}
// Known отвечает, знаком ли слой.
func Known(l Layer) bool { return Rank(l) < len(Layers) }
// Rank — место слоя в порядке от мелкого к крупному. Незнакомый слой считается
// крупнее любого известного.
func Rank(l Layer) int {
for i, known := range Layers {
if known == l {
return i
}
}
return len(Layers)
}
+34
View File
@@ -786,3 +786,37 @@ func TestParseЧужаяПричинаОбрезается(t *testing.T) {
t.Errorf("текст ошибки %d Б: литерал из тела доехал до сообщения", len(err.Error()))
}
}
// Словарь слоёв ОДИН, и порядок в нём — от самого мелкого к самому крупному.
//
// Утверждается прямо, потому что из этого словаря выводятся четыре вещи: ранг
// слоя при разборе, перечень слоёв для выборки охватов при чтении, проверка
// параметра запроса и текст отказа клиенту. Разъехавшись, они дали бы пустой
// ряд при непустых данных — и ни компилятор, ни другой тест этого не увидели бы.
func TestСловарьСлоёвУпорядоченОтМелкогоККрупному(t *testing.T) {
want := []hae.Layer{hae.LayerSample, hae.LayerRaw, hae.LayerMinute, hae.LayerHour, hae.LayerDay}
if len(hae.Layers) != len(want) {
t.Fatalf("слоёв в словаре %d, ожидалось %d", len(hae.Layers), len(want))
}
for i, l := range want {
if hae.Layers[i] != l {
t.Errorf("слой %d — %q, ожидался %q", i, hae.Layers[i], l)
}
if hae.Rank(l) != i {
t.Errorf("ранг %q = %d, ожидался %d", l, hae.Rank(l), i)
}
if !hae.Known(l) {
t.Errorf("слой %q словарём не признан", l)
}
}
// Незнакомый слой крупнее любого известного и словарём не признан: иначе он
// выиграл бы предпочтение при равном охвате.
if got := hae.Rank("weekly"); got != len(hae.Layers) {
t.Errorf("ранг незнакомого слоя %d, ожидался %d", got, len(hae.Layers))
}
if hae.Known("weekly") {
t.Error("незнакомый слой признан словарём")
}
}
+4 -1
View File
@@ -159,5 +159,8 @@ func (a *api) handleMetrics(w http.ResponseWriter, r *http.Request) {
// ответ уходит без метки: это ровно поведение до появления условного
// запроса, то есть деградация в безопасную сторону.
setReadHeaders(w, etag(scopeMetrics, snap.Version))
writeJSON(w, http.StatusOK, catalogWire(snap.Metrics))
// Обрыв записи каталога глушится: ответ здесь — три десятка строк, и
// оборваться на нём нечему. У маршрута точек ряд ничем не ограничен, и там
// отказ записи наблюдается отдельным чекпоинтом.
_ = writeJSON(w, http.StatusOK, catalogWire(snap.Metrics))
}
+30 -3
View File
@@ -16,12 +16,14 @@ import (
"git.vakhrushev.me/av/healthlog/internal/catalog"
"git.vakhrushev.me/av/healthlog/internal/ingest"
"git.vakhrushev.me/av/healthlog/internal/points"
)
// Options — зависимости и настройки транспорта.
type Options struct {
Ingest *ingest.Service
Catalog *catalog.Service
Points *points.Service
Log *slog.Logger
WriteTokens []string
ReadTokens []string
@@ -35,6 +37,7 @@ type Options struct {
type api struct {
ingest *ingest.Service
catalog *catalog.Service
points *points.Service
log *slog.Logger
writeTokens []string
readTokens []string
@@ -47,6 +50,7 @@ func New(o Options) http.Handler {
a := &api{
ingest: o.Ingest,
catalog: o.Catalog,
points: o.Points,
log: o.Log,
writeTokens: o.WriteTokens,
readTokens: o.ReadTokens,
@@ -62,6 +66,10 @@ func New(o Options) http.Handler {
r.Route("/api/v1", func(r chi.Router) {
r.With(requireToken(a.writeTokens)).Post("/ingest", a.handleIngest)
r.With(requireToken(a.readTokens)).Get("/metrics", a.handleMetrics)
// Маршрут точек стоит РЯДОМ с каталогом, а не поверх него: у chi
// литеральный `/metrics` и шаблон `/metrics/{metric}` — разные узлы, и
// каталог остаётся достижим. Утверждается это тестом, а не верой.
r.With(requireToken(a.readTokens)).Get("/metrics/{metric}", a.handlePoints)
})
return r
}
@@ -152,10 +160,27 @@ func routePattern(r *http.Request) string {
return r.URL.Path
}
func writeJSON(w http.ResponseWriter, status int, v any) {
// writeJSON — единственный сериализатор тел ответа.
//
// Экранирование HTML ВЫКЛЮЧЕНО, и это не косметика. `encoding/json` по
// умолчанию превращает `&`, `<` и `>` в `\u0026`, `\u003c`, `\u003e`; на
// маршруте, отдающем дословно сохранённое содержимое точки, это прямо ломает
// обещание дословности — имя источника приходит с телефона пользовательской
// строкой и законно содержит `&`. Хранилище этот же капкан уже проходило и
// обезвредило тем же способом (`store.encodePayload`).
//
// Правило общее для всех читающих маршрутов намеренно: механизм один, и
// решать его заново на каждом маршруте значило бы завести второй способ.
// Отказ записи ВОЗВРАЩАЕТСЯ, а не глушится: код ответа отдан до сериализации,
// поэтому оборванное на середине тело снаружи неотличимо от успеха, а
// `accessLog` честно напишет `200`. Кто из вызывающих обязан об этом сказать —
// решает он сам; глушить молча нельзя ни одному.
func writeJSON(w http.ResponseWriter, status int, v any) error {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(v)
enc := json.NewEncoder(w)
enc.SetEscapeHTML(false)
return enc.Encode(v)
}
// errorWire — форма провода тела отказа, общая для всех маршрутов.
@@ -173,5 +198,7 @@ type errorWire struct {
// writeError отдаёт человекочитаемое сообщение, а не текст ошибки: в тексте
// имена колонок и форма запроса.
func writeError(w http.ResponseWriter, status int, msg string) {
writeJSON(w, status, errorWire{Error: msg})
// Тело отказа — десятки байт: оборваться на нём нечему, и сообщать о
// таком обрыве было бы шумом.
_ = writeJSON(w, status, errorWire{Error: msg})
}
+2
View File
@@ -18,6 +18,7 @@ import (
"git.vakhrushev.me/av/healthlog/internal/catalog"
"git.vakhrushev.me/av/healthlog/internal/httpapi"
"git.vakhrushev.me/av/healthlog/internal/ingest"
"git.vakhrushev.me/av/healthlog/internal/points"
"git.vakhrushev.me/av/healthlog/internal/store"
)
@@ -324,6 +325,7 @@ func newAPILogged(t *testing.T, writeTokens, readTokens []string) (http.Handler,
h := httpapi.New(httpapi.Options{
Ingest: ingest.New(arch, st, nil, log),
Catalog: cat,
Points: points.New(st, log),
Log: log,
WriteTokens: writeTokens,
ReadTokens: readTokens,
+3 -1
View File
@@ -42,7 +42,9 @@ func (a *api) handleIngest(w http.ResponseWriter, r *http.Request) {
return
}
writeJSON(w, http.StatusOK, ingestResponse{
// Ответ приёма — три поля; обрыв записи на нём означает ушедшего клиента, а
// доставка уже сохранена, и это единственное, что здесь обещано.
_ = writeJSON(w, http.StatusOK, ingestResponse{
DeliveryID: res.DeliveryID,
Bytes: res.Bytes,
SHA256: res.SHA256,
+316
View File
@@ -0,0 +1,316 @@
package httpapi
import (
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"net/http"
"net/url"
"strings"
"time"
"github.com/go-chi/chi/v5"
"git.vakhrushev.me/av/healthlog/internal/catalog"
"git.vakhrushev.me/av/healthlog/internal/hae"
"git.vakhrushev.me/av/healthlog/internal/points"
)
// ФОРМА ПРОВОДА ТОЧЕК. Объявлена здесь и только здесь — доменные типы
// `internal/points` `json`-тегов не несут и до сериализации не доезжают.
// Образец и его цена — `internal/httpapi/catalog.go`; выбирать заново не нужно.
// pointsResponse — оболочка ответа точек.
//
// Поля присутствуют ВСЕГДА, даже когда сообщить нечего: клиент не должен
// выводить исход наличием или отсутствием поля. Отсюда указатели у `layer` и
// `bucket` — `null` означает «слоя нет» и «свёртки не было», а не пустую
// строку, которая была бы законным именем слоя.
type pointsResponse struct {
Metric string `json:"metric"`
From time.Time `json:"from"`
To time.Time `json:"to"`
Layer *string `json:"layer"`
// Bucket — сетка свёртки. Всегда `null` в этой версии маршрута: свёртку
// делает соседняя задача. Поле объявлено сразу, потому что менять форму
// конверта после того, как его скопировали четыре маршрута, — не развилка,
// а археология.
Bucket *string `json:"bucket"`
Aggregation seriesAggregation `json:"aggregation"`
Points []pointWire `json:"points"`
}
// seriesAggregation — род свёртки вместе с его применимостью и свежестью.
//
// `applicable` не украшение: род — свойство МЕТРИКИ, слой — свойство ОТДАННОГО
// РЯДА, и сочетание `{"layer": "raw", "style": "cumulative"}` законно, штатно и
// прямо приглашает потребителя сложить интерполяцию самому. Система при этом
// ничего не складывает, а решение у потребителя уже принято по завышенному
// втрое числу.
//
// `last_hour` — ярлык самого свежего часа окна измерения. Окно считается в
// ОБЩИХ часах, а не в часах календаря: выключенная минутная автоматизация HAE
// останавливает их пополнение, окно замирает и продолжает объявлять род.
// Единственный след этого — вот это поле.
type seriesAggregation struct {
Style string `json:"style"`
Applicable bool `json:"applicable"`
LastHour *time.Time `json:"last_hour"`
}
// pointWire — точка ряда на проводе.
//
// `values` уезжает СЫРЫМ JSON: точки хранятся дословно, и переписывать их в тип
// провода значило бы нарушить инвариант. Конверт вокруг значения при этом
// объявлен и нормализован, как того требует «форма Apple не транслируется».
//
// `ts_end` есть потому, что идентичность точки — интервал, а не метка: под
// одной меткой лежит до трёх записей сна. У точки-измерения он равен `ts`.
type pointWire struct {
TS time.Time `json:"ts"`
TSEnd time.Time `json:"ts_end"`
TZOffset int `json:"tz_offset"`
Units string `json:"units"`
Values json.RawMessage `json:"values"`
}
// pointsWire переводит ряд домена в форму провода.
//
// Чистая функция от уже прочитанного значения: ни `context`, ни хранилища, ни
// часов. Версия снимка сюда не идёт — она уезжает в `ETag`.
func pointsWire(s points.Series) pointsResponse {
out := pointsResponse{
Metric: s.Metric,
From: s.From,
To: s.To,
Aggregation: seriesAggregation{
Style: s.Style.String(),
Applicable: s.Applicable,
// Указатель переносится КАК УКАЗАТЕЛЬ: разыменование дало бы
// `0001-01-01T00:00:00Z` там, где окна не было, а правдоподобная
// дата в ответе неотличима от настоящей.
LastHour: s.LastHour,
},
// Непустой срез, а не nil: nil сериализуется в `null`, и клиент
// прочитал бы «поля нет» вместо «точек нет».
Points: make([]pointWire, 0, len(s.Points)),
}
if s.Layer != "" {
layer := s.Layer
out.Layer = &layer
}
for _, p := range s.Points {
out.Points = append(out.Points, pointWire{
TS: p.TS,
TSEnd: p.End,
TZOffset: p.OffsetSeconds,
Units: p.Units,
Values: p.Raw,
})
}
return out
}
// pointsScope — область действия метки ответа: ХЕШ канонизированной формы
// запроса.
//
// Живёт В ТРАНСПОРТЕ, и это решение, а не случайность: у каталога область тоже
// транспортная (`scopeMetrics`), а помощник `etag` прямо просит маршрут положить
// сюда канонизированную форму запроса. Считай её домен — у одного понятия
// оказалось бы два дома, и три следующих маршрута выбирали бы между ними
// монетой. Заодно домен перестал бы носить транспортный артефакт: префикс и hex
// — это форма метки HTTP, а форму провода объявляет транспорт (ADR).
//
// Метка действительна в пределах ОДНОГО набора данных, а ответ этого маршрута
// есть функция параметров. Метка без них однажды подтвердила бы неизменность
// чужого набора — молча и без следов.
//
// Канонизация: границы приводятся к UTC в RFC 3339 с ПОЛНОЙ точностью. Полной,
// а не посекундной, потому что отбор точек идёт по полной метке: два запроса,
// различающиеся долей секунды, дают разные ряды, и одна область на них
// означала бы `304` на чужом наборе данных. Зона при этом канонизируется —
// одно и то же время, записанное разными смещениями, даёт одну область.
//
// Слой берётся ЗАПРОШЕННЫЙ, а не выбранный: областью является запрос, а смену
// выбранного слоя от новых данных ловит версия витрины.
//
// ХЕШ, а не сама форма, и обе причины из `docs/security.md`. Имя метрики
// приходит из чужого тела дословно и ничем не ограничено — в заголовке ответа
// оно дало бы `ETag` в тысячи байт, притом что тот же вход в лог уезжает
// обрезанным. И оно законно содержит кавычку, которая по RFC 9110 кончает
// метку: собственный `scanETag` проекта обрезал бы её ровно там, и условный
// запрос по такой метрике не сработал бы никогда, а симптом увёл бы отладку в
// помощника. Длина впереди остаётся внутри хешируемой строки: без неё имя с
// разделителем склеилось бы с соседним полем.
func pointsScope(req points.Request) string {
canonical := fmt.Sprintf("%d:%s|%s|%s|%s",
len(req.Metric), req.Metric,
req.From.UTC().Format(time.RFC3339Nano),
req.To.UTC().Format(time.RFC3339Nano),
req.Layer)
sum := sha256.Sum256([]byte(canonical))
return "points." + hex.EncodeToString(sum[:scopeBytes])
}
// scopeBytes — сколько байтов хеша попадает в область.
//
// Шестнадцать: 128 бит, то есть столкновение неотличимо от невозможного, а
// заголовок остаётся короче метки каталога. Хеш здесь не криптографический
// секрет — он ограничитель длины и экранирование разом.
const scopeBytes = 16
// handlePoints отдаёт точки метрики за период.
func (a *api) handlePoints(w http.ResponseWriter, r *http.Request) {
req, msg := parsePointsRequest(r)
if msg != "" {
// Отказ до чтения витрины: невозможный запрос не должен стоить снимка.
// В сообщении нет ни одного значения из запроса — оно уедет клиенту, а
// в запросе имя метрики и границы периода.
writeError(w, http.StatusBadRequest, msg)
return
}
series, err := a.points.Series(r.Context(), req)
if err != nil {
// Исход операции логирует доменный слой, транспорт только переводит его
// в ответ. Наружу уходит человекочитаемое сообщение, а не текст ошибки:
// в нём имена колонок и форма запроса.
writeError(w, http.StatusInternalServerError, "ряд точек не собрался")
return
}
// Область действия метки — канонизированная форма запроса: ответ этого
// маршрута есть функция параметров, и метка без них однажды подтвердила бы
// неизменность чужого набора данных. Горизонт измерения в метке уже есть —
// его кладёт домен вместе с версией витрины.
setReadHeaders(w, etag(pointsScope(req), series.Version))
if err := writeJSON(w, http.StatusOK, pointsWire(series)); err != nil {
// Тело оборвалось на середине: дедлайн записи, ушедший клиент, полный
// буфер посредника. Код ответа уже отдан, и `accessLog` напишет `200` —
// то есть единственный канал наблюдаемости сообщит успех о неотданном
// ответе. Ряд ничем не ограничен по размеру (предел — соседняя задача),
// поэтому случай не гипотетический: неделя нижнего слоя это 280 МиБ.
//
// WARN, а не DEBUG, и это исправление находки эксплуатационного прохода:
// боевой уровень логирования — `info` (`config.docker.toml`), то есть
// запись уровня `DEBUG` не прошла бы фильтр НИКОГДА, и единственный
// признак недоставленного тела не существовал бы в проде вовсе.
// Обесценивания уровня здесь нет: событие не периодическое — оно
// означает, что потребитель получил битый JSON.
//
// Значений точек и границ запроса в записи нет; имя метрики обрезано.
a.log.WarnContext(r.Context(), "points response truncated",
"capability", "query", "metric", catalog.ClipMetric(req.Metric), "error", err)
}
}
// parsePointsRequest разбирает параметры маршрута точек.
//
// Второй возврат — человекочитаемая причина отказа; пустая строка означает, что
// запрос принят. Строка, а не ошибка: она целиком уезжает клиенту, поэтому
// обязана быть составлена здесь и не содержать ни одного значения из запроса.
func parsePointsRequest(r *http.Request) (points.Request, string) {
q := r.URL.Query()
// Свёртка по сетке ещё не поддержана, и молчать об этом нельзя: клиент,
// попросивший суточную сетку и получивший полный минутный ряд, заметил бы
// подмену, только пересчитав точки. Это зеркало того самого промаха, ради
// которого «сетка задана явно» вообще различается.
if q.Has("bucket") {
return points.Request{}, "свёртка по сетке ещё не поддержана: параметр bucket не принимается"
}
from, msg := parseBound(q, "from")
if msg != "" {
return points.Request{}, msg
}
to, msg := parseBound(q, "to")
if msg != "" {
return points.Request{}, msg
}
if !from.Before(to) {
return points.Request{}, "период пуст: from должен быть строго раньше to"
}
// `Has`, а не `Get() != ""`, и симметрично `bucket`. Пустое значение
// (`?layer=`) — самый частый способ промахнуться: шаблон клиента с
// невыставленной переменной. Разбор по непустоте молча включил бы
// автоматический выбор, и клиент, спросивший разрез поимённо, не отличил бы
// свой промах от ответа по существу.
layer := q.Get("layer")
if q.Has("layer") && !hae.Known(hae.Layer(layer)) {
return points.Request{}, "неизвестный слой: допустимы " + knownLayers
}
return points.Request{Metric: metricFromPath(r), From: from, To: to, Layer: layer}, ""
}
// parseBound разбирает границу периода.
//
// Только RFC 3339 и только с явным смещением зоны. Голая дата отвергается не из
// строгости: у неё нет зоны, а вопрос «в какой зоне считать сутки» в проекте
// открыт отдельной задачей. Принять её значило бы выбрать зону за клиента молча.
func parseBound(q url.Values, name string) (time.Time, string) {
raw := q.Get(name)
if raw == "" {
return time.Time{}, "параметр " + name + " обязателен"
}
t, err := time.Parse(time.RFC3339, raw)
if err != nil {
return time.Time{}, "параметр " + name +
": ожидается метка времени RFC 3339 с явным смещением зоны, например 2026-01-01T00:00:00Z"
}
// Год после приведения к UTC обязан остаться четырёхзначным, и это не
// придирка к календарю. Хранилище адресует объекты строкой RFC 3339, а
// сравнение границ идёт лексикографически: `9999-12-31T23:00:00-07:00`
// превращается в `10000-01-01T06:00:00Z`, который меньше любой настоящей
// метки как строка, — и запрос молча отдал бы пустой ряд при непустых
// данных. Отказ здесь честнее пустоты: тот же разряд ловится и на входе
// приёма, но там он унаследован и лечится не тут.
if y := t.UTC().Year(); y < 1 || y > 9999 {
return time.Time{}, "параметр " + name + ": год после приведения к UTC вне диапазона 1–9999"
}
return t.UTC(), ""
}
// knownLayers — перечень допустимых слоёв ДЛЯ СООБЩЕНИЯ КЛИЕНТУ, собранный из
// того же словаря, что и проверка. Текст отказа уезжает наружу и потому обязан
// перечислять ровно то, что принимается: разъехавшись, он врал бы клиенту.
var knownLayers = func() string {
out := make([]string, 0, len(hae.Layers))
for _, l := range hae.Layers {
out = append(out, string(l))
}
return strings.Join(out, ", ")
}()
// metricFromPath достаёт имя метрики из пути.
//
// Решение принимается по `RawPath`, а НЕ по успеху декодирования, и это ровно
// то место, где легко ошибиться. `chi` сопоставляет по сырому пути только когда
// тот непуст, то есть когда `net/url` увидел в пути экранирование; тогда
// `{metric}` приезжает закодированным и его надо декодировать. Когда `RawPath`
// пуст, `chi` отдаёт уже декодированное имя — и второе декодирование испортило
// бы его молча.
//
// Цена ошибки построена и прогнана враждебным проходом: метрика с именем
// `a%41b` (имена приходят из тела доставки дословно) кодируется клиентом в
// `a%2541b`, `net/url` декодирует это в `a%41b` и оставляет `RawPath` пустым —
// а второе декодирование давало `aAb`, то есть маршрут отвечал `200` и данными
// ДРУГОЙ метрики.
//
// Отказ декодирования не является отказом маршрута: берём имя как есть — оно
// всё равно не совпадёт ни с одной метрикой витрины, и клиент получит честный
// пустой ряд, а не `400` на своё же имя.
func metricFromPath(r *http.Request) string {
raw := chi.URLParam(r, "metric")
if r.URL.RawPath == "" {
return raw
}
if decoded, err := url.PathUnescape(raw); err == nil {
return decoded
}
return raw
}
+677
View File
@@ -0,0 +1,677 @@
package httpapi_test
import (
"context"
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"net/url"
"strings"
"testing"
"time"
"git.vakhrushev.me/av/healthlog/internal/hae"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// pointsAt — фиксированное «когда» витрины тестов: 1 июня 2026 года.
//
// Час выбран В ПРОШЛОМ и жёстко, а не относительно `time.Now()`, потому что
// окно измерения ограничено горизонтом «сейчас плюс час»: фикстура,
// построенная от текущего времени, при переходе через границу часа дала бы
// другой набор пригодных часов — и другое тело ответа. Класс уже ловили
// (docs/review.md, запись 2026-08-03).
func pointsAt(hh, mm int) time.Time {
return time.Date(2026, 6, 1, hh, mm, 0, 0, time.UTC)
}
func mergePoints(t *testing.T, st *store.Store, in ...store.IncomingPoint) {
t.Helper()
if _, err := st.Merge(context.Background(), store.Incoming{Points: in}, store.DeliveryRef{ID: "d"}); err != nil {
t.Fatalf("слияние: %v", err)
}
}
func pointAt(metric, layer, units string, start, end time.Time, raw string) store.IncomingPoint {
return store.IncomingPoint{
Metric: metric, Layer: layer, Units: units,
Point: store.Point{Start: start, End: end, Raw: json.RawMessage(raw)},
}
}
func getPoints(t *testing.T, h http.Handler, metric, query, auth string) *httptest.ResponseRecorder {
t.Helper()
req := httptest.NewRequest(http.MethodGet, "/api/v1/metrics/"+url.PathEscape(metric)+"?"+query, nil)
if auth != "" {
req.Header.Set("Authorization", auth)
}
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
return rec
}
const dayWindow = "from=2026-06-01T00:00:00Z&to=2026-06-02T00:00:00Z"
// ФОРМА ОТВЕТА ЦЕЛИКОМ, БАЙТАМИ. Детектор изменения публичного контракта: он
// краснеет в момент правки формы, а не у потребителя.
//
// Ни одно поле литерала не зависит от хода часов: витрина фиксирована, часы
// лежат в прошлом, окно измерения на этих данных пусто. Утверждается заодно
// К2 — `layer`, `aggregation` и `last_hour` присутствуют, — и правило «пустая
// коллекция это `[]`, отсутствующее значение это `null`».
func TestТочкиФормаОтветаБайтами(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
mergePoints(t, st,
pointAt("body_mass", "raw", "kg", pointsAt(9, 30), pointsAt(9, 30), `{"qty":72.5,"date":"2026-06-01 12:30:00 +0300"}`),
)
want := `{"metric":"body_mass","from":"2026-06-01T00:00:00Z","to":"2026-06-02T00:00:00Z",` +
`"layer":"raw","bucket":null,` +
`"aggregation":{"style":"unknown","applicable":false,"last_hour":null},` +
`"points":[{"ts":"2026-06-01T09:30:00Z","ts_end":"2026-06-01T09:30:00Z","tz_offset":0,` +
`"units":"kg","values":{"qty":72.5,"date":"2026-06-01 12:30:00 +0300"}}]}`
rec := getPoints(t, h, "body_mass", dayWindow, "")
if rec.Code != http.StatusOK {
t.Fatalf("статус %d, тело %s", rec.Code, rec.Body.String())
}
if got := strings.TrimSpace(rec.Body.String()); got != want {
t.Errorf("форма ответа точек изменилась:\n получили %s\n ждали %s", got, want)
}
}
// К3: род не измерен — сказано словом, а не молчанием, и свёртка не предложена.
// Проверяется через разбор, а не подстрокой: подстрока `"style"` осталась бы на
// месте и у ответа, объявившего род.
func TestТочкиНеизмеренныйРодНазванСловом(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
mergePoints(t, st,
pointAt("six_minute_walking_test_distance", "minute", "m", pointsAt(9, 0), pointsAt(9, 0), `{"qty":420}`),
)
var got struct {
Aggregation struct {
Style string `json:"style"`
Applicable bool `json:"applicable"`
LastHour *string `json:"last_hour"`
} `json:"aggregation"`
Points []json.RawMessage `json:"points"`
}
decodeBody(t, getPoints(t, h, "six_minute_walking_test_distance", dayWindow, ""), &got)
if got.Aggregation.Style != "unknown" {
t.Errorf("род %q, ожидался unknown", got.Aggregation.Style)
}
if got.Aggregation.Applicable {
t.Error("неизмеренный род объявлен применимым — потребителю предложено свернуть по неизвестному")
}
if got.Aggregation.LastHour != nil {
t.Errorf("граница окна %v, ожидался null", *got.Aggregation.LastHour)
}
if len(got.Points) != 1 {
t.Errorf("точек %d, ожидалась 1: род неизвестен — точки отдаются как есть", len(got.Points))
}
}
// Измеренный род едет вместе с границей окна и с применимостью. Данные —
// минутный и часовой слои одной метрики, сходящиеся суммой: ровно тот вход, на
// котором каталог объявляет `cumulative`.
func TestТочкиИзмеренныйРодЕдетСГраницейОкна(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
var in []store.IncomingPoint
for hh := range catalogAgreeingHours {
hour := pointsAt(hh, 0)
// Часовое значение равно сумме минутных, и сумма отличима от среднего.
in = append(in,
pointAt("active_energy", "hour", "kJ", hour, hour, `{"qty":30}`),
pointAt("active_energy", "minute", "kJ", hour, hour, `{"qty":10}`),
pointAt("active_energy", "minute", "kJ", hour.Add(time.Minute), hour.Add(time.Minute), `{"qty":20}`),
)
}
mergePoints(t, st, in...)
cases := []struct {
name string
query string
wantLayer string
wantApplicable bool
}{
{"часовой слой", dayWindow + "&layer=hour", "hour", true},
{"минутный слой", dayWindow + "&layer=minute", "minute", true},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
var got struct {
Layer *string `json:"layer"`
Aggregation struct {
Style string `json:"style"`
Applicable bool `json:"applicable"`
LastHour *string `json:"last_hour"`
} `json:"aggregation"`
}
decodeBody(t, getPoints(t, h, "active_energy", c.query, ""), &got)
if got.Aggregation.Style != "cumulative" {
t.Fatalf("род %q, ожидался cumulative", got.Aggregation.Style)
}
if got.Layer == nil || *got.Layer != c.wantLayer {
t.Errorf("слой %v, ожидался %q", got.Layer, c.wantLayer)
}
if got.Aggregation.Applicable != c.wantApplicable {
t.Errorf("применимость %v, ожидалась %v", got.Aggregation.Applicable, c.wantApplicable)
}
// Граница окна — ярлык часа, а не конец периода и не «сейчас».
if got.Aggregation.LastHour == nil {
t.Fatal("граница окна измерения null при измеренном роде")
}
if want := "2026-06-01T13:00:00Z"; *got.Aggregation.LastHour != want {
t.Errorf("граница окна %q, ожидалась %q", *got.Aggregation.LastHour, want)
}
})
}
}
// catalogAgreeingHours — сколько согласных часов кладётся в фикстуру. Больше
// порога каталога, и число названо здесь, а не подобрано в теле теста.
const catalogAgreeingHours = 14
// Накопительная метрика на нижнем слое HAE объявляется НЕПРИМЕНИМОЙ: нижний
// слой это интерполяция, а не сэмплы, и сумма по нему завышает втрое. Без этого
// поля конверт приглашает агента сложить её самому.
func TestТочкиНакопительнаяМетрикаНаНижнемСлоеНеприменима(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
var in []store.IncomingPoint
for hh := range catalogAgreeingHours {
hour := pointsAt(hh, 0)
in = append(in,
pointAt("active_energy", "hour", "kJ", hour, hour, `{"qty":30}`),
pointAt("active_energy", "minute", "kJ", hour, hour, `{"qty":10}`),
pointAt("active_energy", "minute", "kJ", hour.Add(time.Minute), hour.Add(time.Minute), `{"qty":20}`),
pointAt("active_energy", "raw", "kJ", hour.Add(7*time.Second), hour.Add(7*time.Second), `{"qty":1}`),
)
}
mergePoints(t, st, in...)
var got struct {
Aggregation struct {
Style string `json:"style"`
Applicable bool `json:"applicable"`
} `json:"aggregation"`
}
decodeBody(t, getPoints(t, h, "active_energy", dayWindow+"&layer=raw", ""), &got)
if got.Aggregation.Style != "cumulative" {
t.Fatalf("род %q, ожидался cumulative", got.Aggregation.Style)
}
if got.Aggregation.Applicable {
t.Error("накопительная метрика на нижнем слое объявлена применимой — потребитель просуммирует интерполяцию")
}
}
// Содержимое точки уезжает ДОСЛОВНО, включая символы, которые `encoding/json`
// по умолчанию превращает в escape-последовательности.
//
// Случай заведён отдельно и намеренно: фикстуры `testdata` символов `&<>` не
// содержат вовсе, то есть утверждение о дословности на них зелено и будучи
// сломанным. Имя источника приходит с телефона пользовательской строкой.
func TestТочкиСодержимоеНеЭкранируется(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
raw := `{"qty":1,"source":"iPhone <A&B>"}`
mergePoints(t, st, pointAt("m", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), raw))
body := getPoints(t, h, "m", dayWindow, "").Body.String()
if !strings.Contains(body, raw) {
t.Errorf("содержимое точки переписано сериализатором:\n %s", body)
}
for _, escaped := range []string{`\u0026`, `\u003c`, `\u003e`} {
if strings.Contains(body, escaped) {
t.Errorf("сериализатор заэкранировал содержимое точки (%s) — обещание дословности нарушено", escaped)
}
}
}
// Интервальная точка отдаёт КОНЕЦ координаты: под одной меткой лежит до трёх
// записей сна, и конверт с одним `ts` предложил бы клиенту различать их,
// разбирая дословное содержимое.
func TestТочкиИнтервалОтдаётКонецКоординаты(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
start, end := pointsAt(1, 0), pointsAt(3, 30)
mergePoints(t, st,
pointAt("sleep_analysis", "raw", "hr", start, end, `{"value":"В кровати"}`),
pointAt("sleep_analysis", "raw", "hr", start, pointsAt(2, 0), `{"value":"Глубокий"}`),
)
var got struct {
Points []struct {
TS string `json:"ts"`
TSEnd string `json:"ts_end"`
} `json:"points"`
}
decodeBody(t, getPoints(t, h, "sleep_analysis", dayWindow, ""), &got)
if len(got.Points) != 2 {
t.Fatalf("точек %d, ожидалось 2: одна метка, разные интервалы", len(got.Points))
}
if got.Points[0].TS != got.Points[1].TS {
t.Fatalf("метки разошлись: %s и %s", got.Points[0].TS, got.Points[1].TS)
}
if got.Points[0].TSEnd == got.Points[1].TSEnd {
t.Errorf("концы координат совпали (%s) — записи неразличимы", got.Points[0].TSEnd)
}
// Порядок детерминирован: при равном начале раньше идёт более короткий.
if got.Points[0].TSEnd != "2026-06-01T02:00:00Z" {
t.Errorf("порядок точек не по (ts, ts_end): первый конец %s", got.Points[0].TSEnd)
}
}
// Пустой период — успех, а не отсутствие ресурса. `layer` при этом `null`
// ТОЛЬКО когда слой выбирала система: клиент, спросивший разрез поимённо,
// обязан отличать «этого разреза за период нет» от «параметр проигнорирован».
func TestТочкиПустойПериодОтличаетВыбранныйСлойОтЗапрошенного(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
mergePoints(t, st, pointAt("m", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":1}`))
empty := "from=2026-07-01T00:00:00Z&to=2026-07-02T00:00:00Z"
cases := []struct {
name string
query string
want string
}{
{"слой выбирала система", empty, "null"},
{"слой задан явно", empty + "&layer=minute", `"minute"`},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
var got struct {
Layer json.RawMessage `json:"layer"`
Points []json.RawMessage `json:"points"`
}
rec := getPoints(t, h, "m", c.query, "")
if rec.Code != http.StatusOK {
t.Fatalf("статус %d, ожидался 200: «данных нет» — не «ресурса нет»", rec.Code)
}
decodeBody(t, rec, &got)
if string(got.Layer) != c.want {
t.Errorf("слой %s, ожидался %s", got.Layer, c.want)
}
if got.Points == nil {
t.Error("точки уехали как null — клиент прочитает «поля нет» вместо «точек нет»")
}
})
}
}
// Имени метрики в витрине нет вовсе — это тоже `200` с пустым рядом: список
// имён маршруту точек не принадлежит, их отдаёт каталог.
func TestТочкиНезнакомойМетрикиОтвечают200(t *testing.T) {
h, _, _ := newAPITokens(t, nil, nil)
rec := getPoints(t, h, "нет такой метрики", dayWindow, "")
if rec.Code != http.StatusOK {
t.Fatalf("статус %d, ожидался 200", rec.Code)
}
}
// Имя метрики достаётся из пути ДЕКОДИРОВАННЫМ. Без этого метрика с пробелом
// или слэшем в имени была бы недостижима, а имена приходят из тела доставки
// дословно и ничем не ограничены.
func TestТочкиИмяМетрикиДекодируетсяИзПути(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
for _, metric := range []string{"с пробелом", "со/слэшем", "с%знаком"} {
t.Run(metric, func(t *testing.T) {
mergePoints(t, st, pointAt(metric, "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":1}`))
var got struct {
Metric string `json:"metric"`
Points []json.RawMessage `json:"points"`
}
decodeBody(t, getPoints(t, h, metric, dayWindow, ""), &got)
if got.Metric != metric {
t.Errorf("имя метрики %q, ожидалось %q", got.Metric, metric)
}
if len(got.Points) != 1 {
t.Errorf("точек %d, ожидалась 1 — метрика недостижима по своему имени", len(got.Points))
}
})
}
}
// Маршрут точек НЕ перехватывает каталог: у chi литеральный `/metrics` и
// шаблон `/metrics/{metric}` — разные узлы, но верить в это нельзя.
func TestТочкиНеПерехватываютКаталог(t *testing.T) {
h, _, _ := newAPITokens(t, nil, nil)
if got := strings.TrimSpace(getCatalog(t, h, "").Body.String()); got != `{"metrics":[]}` {
t.Errorf("каталог перехвачен маршрутом точек: %s", got)
}
}
// Разбор параметров: каждый отказ — `400`, до чтения витрины, и без единого
// значения из запроса в теле ответа.
func TestТочкиОтвергаютНевозможныйЗапрос(t *testing.T) {
h, _, _ := newAPITokens(t, nil, nil)
cases := map[string]string{
"нет from": "to=2026-06-02T00:00:00Z",
"нет to": "from=2026-06-01T00:00:00Z",
"голая дата": "from=2026-06-01&to=2026-06-02",
"без зоны": "from=2026-06-01T00:00:00&to=2026-06-02T00:00:00Z",
"мусор": "from=вчера&to=сегодня",
"вывернутый период": "from=2026-06-02T00:00:00Z&to=2026-06-01T00:00:00Z",
"пустой период": "from=2026-06-01T00:00:00Z&to=2026-06-01T00:00:00Z",
"незнакомый слой": dayWindow + "&layer=weekly",
"свёртка не поддержана": dayWindow + "&bucket=day",
}
for name, query := range cases {
t.Run(name, func(t *testing.T) {
rec := getPoints(t, h, "секретное_имя_метрики", query, "")
if rec.Code != http.StatusBadRequest {
t.Fatalf("статус %d, ожидался 400 (тело %s)", rec.Code, rec.Body.String())
}
// Значения из запроса в тело отказа не уезжают: там имя метрики и
// границы периода, а тело отказа читает кто угодно.
for _, leak := range []string{"2026-06-01", "2026-06-02", "weekly", "вчера", "секретное_имя_метрики"} {
if strings.Contains(rec.Body.String(), leak) {
t.Errorf("в теле отказа значение из запроса (%q): %s", leak, rec.Body.String())
}
}
})
}
}
// Незнакомые параметры игнорируются, как принято в HTTP: клиент, приславший
// лишнее, получает данные, а не отказ.
func TestТочкиИгнорируютНезнакомыйПараметр(t *testing.T) {
h, _, _ := newAPITokens(t, nil, nil)
rec := getPoints(t, h, "m", dayWindow+"&limit=10", "")
if rec.Code != http.StatusOK {
t.Fatalf("статус %d, ожидался 200", rec.Code)
}
}
func TestТочкиТребуютТокенЧтения(t *testing.T) {
h, _, _ := newAPITokens(t, []string{"write-token"}, []string{"read-token"})
cases := []struct {
name string
auth string
want int
}{
{"без заголовка", "", http.StatusUnauthorized},
{"токен приёма", "Bearer write-token", http.StatusUnauthorized},
{"токен чтения", "Bearer read-token", http.StatusOK},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
if got := getPoints(t, h, "m", dayWindow, c.auth).Code; got != c.want {
t.Errorf("статус %d, ожидался %d", got, c.want)
}
})
}
}
// Метка ответа различает каждый параметр по очереди при ОДНОМ состоянии
// витрины: метка, их не различающая, однажды подтвердит неизменность чужого
// набора данных. Заодно — метка точек отличается от метки каталога.
func TestМеткаТочекРазличаетЗапросы(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
mergePoints(t, st,
pointAt("m", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":1}`),
pointAt("other", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":1}`),
)
base := getPoints(t, h, "m", dayWindow, "").Header().Get("ETag")
if base == "" {
t.Fatal("ответ точек ушёл без метки — условному запросу не на чем стоять")
}
if base == getCatalog(t, h, "").Header().Get("ETag") {
t.Error("метка точек совпала с меткой каталога")
}
others := map[string]func() string{
"метрика": func() string { return getPoints(t, h, "other", dayWindow, "").Header().Get("ETag") },
"начало": func() string {
return getPoints(t, h, "m", "from=2026-06-01T01:00:00Z&to=2026-06-02T00:00:00Z", "").Header().Get("ETag")
},
"конец": func() string {
return getPoints(t, h, "m", "from=2026-06-01T00:00:00Z&to=2026-06-03T00:00:00Z", "").Header().Get("ETag")
},
"слой": func() string { return getPoints(t, h, "m", dayWindow+"&layer=raw", "").Header().Get("ETag") },
}
for name, get := range others {
if got := get(); got == base {
t.Errorf("метка не различает %s: %s", name, got)
}
}
// Эквивалентная запись границ даёт ту же метку: иначе условный запрос не
// сработал бы у клиента, пишущего смещение зоны иначе, чем сервер.
same := getPoints(t, h, "m", "from=2026-06-01T03:00:00%2B03:00&to=2026-06-02T03:00:00%2B03:00", "").Header().Get("ETag")
if same != base {
t.Errorf("эквивалентная запись дала другую метку:\n %s\n %s", same, base)
}
}
// Ответ чтения непригоден для разделяемого кеша: при выключенной проверке
// токенов в запросе нет и `Authorization`, и выгрузку истории здоровья вправе
// подержать у себя любой прокси на пути.
func TestТочкиПомеченыЧастнымКешем(t *testing.T) {
h, _, _ := newAPITokens(t, nil, nil)
if got := getPoints(t, h, "m", dayWindow, "").Header().Get("Cache-Control"); got != "private, no-cache" {
t.Errorf("Cache-Control %q, ожидался private, no-cache", got)
}
}
func decodeBody(t *testing.T, rec *httptest.ResponseRecorder, v any) {
t.Helper()
if rec.Code != http.StatusOK {
t.Fatalf("статус %d, тело %s", rec.Code, rec.Body.String())
}
if err := json.Unmarshal(rec.Body.Bytes(), v); err != nil {
t.Fatalf("разбор ответа: %v (тело %s)", err, rec.Body.String())
}
}
// Отказ хранилища переводится в `500` с человекочитаемым сообщением: наружу не
// уходит ни текст ошибки (в нём имена колонок), ни пустой ряд, который клиент
// прочитал бы как «данных нет».
func TestТочкиОтказХранилищаДаёт500(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
mergePoints(t, st, pointAt("m", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":1}`))
if err := st.Close(); err != nil {
t.Fatalf("закрытие: %v", err)
}
rec := getPoints(t, h, "m", dayWindow, "")
if rec.Code != http.StatusInternalServerError {
t.Fatalf("статус %d, ожидался 500 (тело %s)", rec.Code, rec.Body.String())
}
if strings.Contains(rec.Body.String(), "bucket") || strings.Contains(rec.Body.String(), "select") {
t.Errorf("в теле отказа внутренности хранилища: %s", rec.Body.String())
}
}
// Имя метрики, содержащее процентную последовательность, обязано доехать до
// витрины НЕИЗМЕНЁННЫМ.
//
// Путь построен враждебным проходом и стоил ответа данными ЧУЖОЙ метрики:
// клиент кодирует `a%41b` в `a%2541b`, `net/url` декодирует это обратно в
// `a%41b` и оставляет `RawPath` пустым, а второе декодирование давало `aAb` —
// имя соседней метрики, лежащей рядом в витрине.
func TestТочкиИмяМетрикиНеДекодируетсяДважды(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
// Обе метрики лежат рядом: подмена наблюдаема числом точек.
mergePoints(t, st,
pointAt("a%41b", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":1}`),
pointAt("aAb", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":2}`),
pointAt("aAb", "raw", "count", pointsAt(9, 1), pointsAt(9, 1), `{"qty":3}`),
)
for _, metric := range []string{"a%41b", "aAb", "100%", "a%zzb", "a+b", "шаги", "со/слэшем", "с пробелом"} {
t.Run(metric, func(t *testing.T) {
var got struct {
Metric string `json:"metric"`
}
decodeBody(t, getPoints(t, h, metric, dayWindow, ""), &got)
if got.Metric != metric {
t.Errorf("маршрут ответил о метрике %q, спрашивали %q", got.Metric, metric)
}
})
}
// И прямая проверка исхода: у `a%41b` одна точка, у `aAb` — две.
var got struct {
Points []json.RawMessage `json:"points"`
}
decodeBody(t, getPoints(t, h, "a%41b", dayWindow, ""), &got)
if len(got.Points) != 1 {
t.Errorf("точек %d, ожидалась 1 — ответ собран по чужой метрике", len(got.Points))
}
}
// Пустое значение `?layer=` — промах клиента, а не «слой не задан». Разбор по
// непустоте молча включал бы автоматический выбор, и клиент, спросивший разрез
// поимённо, не отличил бы свой промах от ответа по существу.
func TestТочкиОтвергаютПустойСлой(t *testing.T) {
h, _, _ := newAPITokens(t, nil, nil)
if got := getPoints(t, h, "m", dayWindow+"&layer=", "").Code; got != http.StatusBadRequest {
t.Errorf("статус %d, ожидался 400", got)
}
}
// Границы, различающиеся ДОЛЯМИ СЕКУНДЫ, дают разные метки: отбор точек идёт по
// полной метке, значит и ряды разные. Путь построен враждебным проходом — две
// побайтово одинаковые метки на разных телах, то есть будущий `304` на чужом
// наборе данных.
func TestМеткаТочекРазличаетДолиСекунды(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
mergePoints(t, st, pointAt("m", "raw", "count", pointsAt(0, 0), pointsAt(0, 0), `{"qty":1}`))
whole := getPoints(t, h, "m", "from=2026-06-01T00:00:00Z&to=2026-06-02T00:00:00Z", "")
fraction := getPoints(t, h, "m", "from=2026-06-01T00:00:00.500Z&to=2026-06-02T00:00:00Z", "")
if a, b := whole.Header().Get("ETag"), fraction.Header().Get("ETag"); a == b {
t.Errorf("метка не различает доли секунды: %s", a)
}
if whole.Body.String() == fraction.Body.String() {
t.Error("тела совпали — вход подобран неверно, утверждение о метках ничего не доказывает")
}
}
// Метка ОГРАНИЧЕНА по длине и не выносит наружу имя метрики.
//
// Имя приходит из чужого тела дословно и ничем не ограничено: без предела
// заголовок `ETag` разрастался вместе с ним (прогнано: имя в 3000 байт давало
// заголовок в 3099). Кавычка в имени по RFC 9110 кончает метку, и собственный
// `scanETag` проекта обрезал бы её ровно там — условный запрос по такой метрике
// не сработал бы никогда.
func TestМеткаТочекОграниченаИНеНесётИмя(t *testing.T) {
h, _, _ := newAPITokens(t, nil, nil)
long := strings.Repeat("щ", 1000) + `quote"inside`
tag := getPoints(t, h, long, dayWindow, "").Header().Get("ETag")
if tag == "" {
t.Fatal("ответ ушёл без метки")
}
if len(tag) > 128 {
t.Errorf("метка в %d байт — имя метрики уехало в заголовок целиком", len(tag))
}
if strings.Contains(tag, "щ") || strings.Contains(tag, `quote"`) {
t.Errorf("имя метрики видно в метке: %s", tag)
}
}
// brokenWriter — ResponseWriter, отказывающий на записи тела. Заголовки и код
// он принимает: обрыв случается ПОСЛЕ того, как `200` уже отдан, — ровно так
// это выглядит при сработавшем дедлайне записи или ушедшем клиенте.
type brokenWriter struct {
header http.Header
status int
}
func (w *brokenWriter) Header() http.Header { return w.header }
func (w *brokenWriter) WriteHeader(s int) { w.status = s }
func (w *brokenWriter) Write([]byte) (int, error) {
return 0, errors.New("соединение оборвано")
}
// Обрыв записи тела оставляет собственный чекпоинт.
//
// Код ответа отдан до сериализации, поэтому `accessLog` честно напишет `200` —
// то есть единственный канал наблюдаемости сообщит успех о неотданном ответе.
// Путь построен враждебным проходом на настоящем сокете: тело в 13 МиБ
// оборвалось на 2.7 МиБ, клиент получил нечитаемый JSON, лог — `200`.
func TestТочкиОбрывЗаписиТелаНаблюдаем(t *testing.T) {
h, st, _, seen := newAPILogged(t, nil, nil)
mergePoints(t, st, pointAt("m", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":1}`))
seen.reset()
req := httptest.NewRequest(http.MethodGet, "/api/v1/metrics/m?"+dayWindow, nil)
w := &brokenWriter{header: http.Header{}}
h.ServeHTTP(w, req)
if w.status != http.StatusOK {
t.Fatalf("статус %d, ожидался 200: обрыв случается после кода ответа", w.status)
}
if !seen.has("points response truncated") {
t.Error("обрыв записи прошёл молча — владелец увидит только успешный 200")
}
}
// Словарь слоёв ОДИН: `hae.Layers`. Три копии — порядок, перечень выборки и
// текст отказа — разошлись бы молча, и симптомом был бы пустой ряд при
// непустых данных.
func TestСловарьСлоёвОдин(t *testing.T) {
h, _, _ := newAPITokens(t, nil, nil)
for _, l := range hae.Layers {
if got := getPoints(t, h, "m", dayWindow+"&layer="+string(l), "").Code; got != http.StatusOK {
t.Errorf("слой %q отвергнут статусом %d, хотя он в словаре", l, got)
}
}
rec := getPoints(t, h, "m", dayWindow+"&layer=weekly", "")
if rec.Code != http.StatusBadRequest {
t.Fatalf("статус %d, ожидался 400", rec.Code)
}
for _, l := range hae.Layers {
if !strings.Contains(rec.Body.String(), string(l)) {
t.Errorf("текст отказа не называет слой %q: %s", l, rec.Body.String())
}
}
}
// Граница, уезжающая за четырёхзначный год, отвергается, а не отдаёт пустой ряд.
//
// Объекты адресуются строкой RFC 3339, границы сравниваются лексикографически:
// `9999-12-31T23:00:00-07:00` становится `10000-01-01T06:00:00Z`, который как
// строка меньше любой настоящей метки. Без отказа запрос молча вернул бы пустой
// ряд при непустых данных — найдено триажем.
func TestТочкиОтвергаютГраницуЗаЧетырёхзначнымГодом(t *testing.T) {
h, st, _ := newAPITokens(t, nil, nil)
mergePoints(t, st, pointAt("m", "raw", "count", pointsAt(9, 0), pointsAt(9, 0), `{"qty":1}`))
rec := getPoints(t, h, "m", "from=2020-01-01T00:00:00Z&to=9999-12-31T23:00:00-07:00", "")
if rec.Code != http.StatusBadRequest {
t.Fatalf("статус %d, ожидался 400 (тело %s)", rec.Code, rec.Body.String())
}
}
+1
View File
@@ -101,6 +101,7 @@ func foreignTypes(t reflect.Type) []string {
func TestФормаПроводаДоменаНеСодержит(t *testing.T) {
cases := map[string]any{
"каталог": catalogResponse{},
"точки": pointsResponse{},
"тело отказа": errorWire{},
"учёт приёма": ingestResponse{},
}
+148
View File
@@ -0,0 +1,148 @@
package points_test
import (
"context"
"encoding/json"
"log/slog"
"os"
"path/filepath"
"runtime"
"testing"
"time"
"git.vakhrushev.me/av/healthlog/internal/hae"
"git.vakhrushev.me/av/healthlog/internal/points"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// ЦЕНА МАРШРУТА ТОЧЕК, ЗАМЕР.
//
// Предела размера ответа у маршрута нет намеренно — его вводит соседняя задача
// `read-api-response-limit`. Поэтому цена обязана быть НАЗВАНА ЧИСЛОМ, и число
// снимается на том режиме, ради которого предел заводится, а не на том, который
// оказался под рукой: прецедент 2026-08-04 (`docs/review.md`) — замер, снятый на
// корпусе, где измеряемого случая не бывает, стоил решения «индекс не нужен».
//
// Три режима, от лёгкого к худшему:
//
// редкая метрика за год — «вес за год», сценарий постановки;
// плотная метрика за сутки — минутный слой, штатный запрос трекера;
// плотная метрика за неделю в нижнем слое — то, что правило выбора слоя
// отдаёт по умолчанию, когда охваты равны.
//
// Точки размножаются из РЕАЛЬНОГО пакета HAE (`internal/hae/testdata`), а не
// пишутся литералами: форма точки, длина строк и вид числового литерала входят
// в цену — они определяют и объём gzip, и работу разжатия.
//
// Прогон: go test ./internal/points -run XXX -bench . -benchtime 1x
func BenchmarkРядТочек(b *testing.B) {
cases := []struct {
name string
layer hae.Layer
hours int
perDay int
}{
// Год, взвешивание примерно раз в сутки.
{"редкая метрика за год", hae.LayerHour, 365 * 24, 1},
// Сутки минутного слоя: 60 точек в час.
{"плотная метрика за сутки, minute", hae.LayerMinute, 24, 24 * 60},
// Неделя нижнего слоя. Порядок взят из разведки: около 100 тысяч
// координат в сутки на весь поток; на одну плотную метрику — 3600 в час.
{"плотная метрика за неделю, raw", hae.LayerRaw, 7 * 24, 24 * 3600},
}
for _, c := range cases {
b.Run(c.name, func(b *testing.B) {
svc, from, to := benchStore(b, c.layer, c.hours, c.perDay)
var before, after runtime.MemStats
runtime.GC()
runtime.ReadMemStats(&before)
b.ResetTimer()
var got points.Series
for range b.N {
var err error
got, err = svc.Series(context.Background(), points.Request{
Metric: "bench_metric", From: from, To: to, Layer: string(c.layer),
})
if err != nil {
b.Fatalf("Series: %v", err)
}
}
b.StopTimer()
runtime.ReadMemStats(&after)
b.ReportMetric(float64(len(got.Points)), "точек")
b.ReportMetric(float64(after.TotalAlloc-before.TotalAlloc)/(1<<20)/float64(b.N), "МиБ_выделено")
})
}
}
// benchStore наполняет витрину точками, размноженными из реального пакета HAE.
func benchStore(b *testing.B, layer hae.Layer, hours, perDay int) (*points.Service, time.Time, time.Time) {
b.Helper()
raw := realPointBody(b)
st, err := store.Open(filepath.Join(b.TempDir(), "healthlog.db"))
if err != nil {
b.Fatalf("store.Open: %v", err)
}
b.Cleanup(func() { _ = st.Close() })
perHour := max(perDay/24, 1)
step := time.Hour / time.Duration(perHour)
start := time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC)
ctx := context.Background()
// Пачками по часу: одна транзакция на весь год держала бы блокировку минуты.
for h := range hours {
hour := start.Add(time.Duration(h) * time.Hour)
in := make([]store.IncomingPoint, 0, perHour)
for i := range perHour {
at := hour.Add(time.Duration(i) * step)
in = append(in, store.IncomingPoint{
Metric: "bench_metric", Layer: string(layer), Units: "count",
Point: store.Point{Start: at, End: at, OffsetSeconds: 10800, Raw: raw},
})
}
if _, err := st.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "bench"}); err != nil {
b.Fatalf("слияние: %v", err)
}
}
svc := points.New(st, slog.New(slog.DiscardHandler))
return svc, start, start.Add(time.Duration(hours) * time.Hour)
}
// realPointBody достаёт содержимое настоящей точки из пакета HAE.
//
// Конвенция проекта: тесты формата держим на реальных пакетах. Здесь она нужна
// не ради разбора, а ради ЦЕНЫ — выдуманная точка `{"qty":1}` жмётся иначе и
// разжимается быстрее, чем то, что реально шлёт телефон.
func realPointBody(b *testing.B) json.RawMessage {
b.Helper()
body, err := os.ReadFile(filepath.Join("..", "hae", "testdata", "minute.json"))
if err != nil {
b.Fatalf("чтение пакета: %v", err)
}
var pkg struct {
Data struct {
Metrics []struct {
Data []json.RawMessage `json:"data"`
} `json:"metrics"`
} `json:"data"`
}
if err := json.Unmarshal(body, &pkg); err != nil {
b.Fatalf("разбор пакета: %v", err)
}
for _, m := range pkg.Data.Metrics {
if len(m.Data) > 0 {
return m.Data[0]
}
}
b.Fatal("в пакете нет ни одной точки")
return nil
}
+244
View File
@@ -0,0 +1,244 @@
// Package points — ряд точек одной метрики за период.
//
// Отвечает на вопрос потребителя «дай значения» — в отличие от каталога,
// который отвечает «что у тебя есть». Пакет производен от витрины и ничего в
// неё не пишет.
//
// Три решения этого узла живут здесь, потому что все три — правила, а не
// выборки:
//
// - какой слой отдать, когда клиент его не назвал;
// - применим ли измеренный род к отданному ряду (нижний слой HAE не
// суммируется никогда, а род — свойство метрики, не ряда);
// - из чего собрана метка ответа.
//
// Род при этом НЕ измеряется здесь: правило измерения одно и живёт в
// `internal/catalog`. Второй его экземпляр разошёлся бы с первым молча, а на
// роде строится арифметика года.
package points
import (
"context"
"encoding/json"
"log/slog"
"time"
"git.vakhrushev.me/av/healthlog/internal/catalog"
"git.vakhrushev.me/av/healthlog/internal/hae"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// ФОРМЫ ПРОВОДА В ЭТОМ ПАКЕТЕ НЕТ. Типы ниже — форма ответа use-case:
// `json`-тегов они не несут и до сериализации не доезжают. Публичный контракт
// объявляет транспорт (`internal/httpapi`), см. ADR о форме провода.
// Request — запрос ряда. Границы уже разобраны и нормализованы транспортом;
// период — полуинтервал `[From, To)`.
type Request struct {
Metric string
From time.Time
To time.Time
// Layer — явно запрошенный слой; пустая строка означает «выбери сам».
Layer string
}
// Point — точка ряда.
//
// `Raw` — содержимое ровно в том виде, в каком его сохранило хранилище: точки
// хранятся дословно, и нормализовано у них только время.
type Point struct {
TS time.Time
End time.Time
OffsetSeconds int
Units string
Raw json.RawMessage
}
// Series — ответ маршрута точек вместе с версией ответа.
//
// Версия пустая, когда подписать ответ нечем: витрина изменилась, пока ответ
// собирался, или прочитать её версию не удалось. Это не отказ.
type Series struct {
Version string
Metric string
From time.Time
To time.Time
// Layer — слой, из которого собран ряд. Пустая строка означает, что слой
// выбирала система и выбирать было не из чего; явно запрошенный слой
// уезжает здесь всегда, даже когда ряд пуст.
Layer string
// Style — измеренный род метрики, тем же правилом и тем же окном, что у
// каталога.
Style catalog.Style
// Applicable — применим ли объявленный род к ЭТОМУ ряду.
Applicable bool
// LastHour — ярлык самого свежего часа окна измерения; nil при пустом окне.
LastHour *time.Time
Points []Point
}
// Service собирает ряд точек по витрине.
type Service struct {
store *store.Store
log *slog.Logger
}
// New собирает сервис точек.
func New(st *store.Store, log *slog.Logger) *Service {
return &Service{store: st, log: log}
}
// Series собирает ряд точек метрики за период.
//
// Всё, от чего зависит ответ, снимается ОДНОЙ транзакцией чтения: охваты слоёв,
// точки выбранного слоя и объекты окна измерения. Пара проб версии вокруг неё
// нужна для метки, а непротиворечивость тела даёт транзакция — проба
// расхождение обнаружила бы, но тело всё равно уехало бы клиенту.
//
// Горизонт снимается ОДИН РАЗ и уходит и в окно измерения, и в метку: род есть
// функция горизонта, а горизонт едет вместе с часами. Сними их порознь — и
// метка однажды подтвердит неизменность ответа, чей род уже перевернулся.
func (s *Service) Series(ctx context.Context, req Request) (Series, error) {
horizon := catalog.Horizon()
var snap store.SeriesSnapshot
version, err := s.store.VersionedRead(ctx, func(ctx context.Context) error {
var err error
snap, err = s.store.ReadSeries(ctx, store.SeriesWindow{
Metric: req.Metric,
From: req.From,
To: req.To,
Layer: req.Layer,
Layers: layerNames,
Measure: catalog.MeasureWindow(horizon),
}, func(spans []store.LayerSpan) string {
return pickLayer(req.From, req.To, spans)
})
return err
})
if err != nil { //nolint:nestif // ветка одна, вложенность даёт лог по адресату
// Единственный логирующий чекпоинт исхода: транспорт переводит ошибку в
// ответ и второй раз её не пишет. Границ запроса и значений точек в
// записи нет — данные о здоровье чувствительнее токенов; имя метрики
// обрезано тем же пределом, что у каталога: оно приходит из тела
// дословно, а запись повторяется на каждый запрос.
//
// Отмена снаружи и занятость базы означают «не сделано», а не «не
// выходит»: клиент, оборвавший запрос по своему тайм-ауту, не должен
// давать владельцу ERROR.
if store.Transient(err) {
s.log.DebugContext(ctx, "series interrupted", "capability", "query",
"metric", catalog.ClipMetric(req.Metric), "error", err)
} else {
s.log.ErrorContext(ctx, "series failed", "capability", "query",
"metric", catalog.ClipMetric(req.Metric), "error", err)
}
return Series{}, err
}
// Предупреждения измерения (противоречащий род, данные из будущего) здесь
// НЕ пишутся: они привилегия каталога. Агент опрашивает по расписанию, и
// WARN на каждый опрос обесценил бы уровень ровно так же, как обесценила бы
// его строка на каждый `304`.
style, basis := catalog.Measure(snap.Pairs)
out := Series{
Version: catalog.Stamp(version, horizon),
Metric: req.Metric,
From: req.From,
To: req.To,
Layer: snap.Layer,
Style: style,
Applicable: applicable(style, snap.Layer),
LastHour: basis.LastHour,
// Непустой срез, а не nil: пустой ряд обязан уехать клиенту как `[]`.
Points: make([]Point, 0, len(snap.Points)),
}
for _, p := range snap.Points {
out.Points = append(out.Points, Point{
TS: p.Start,
End: p.End,
OffsetSeconds: p.OffsetSeconds,
Units: p.Units,
Raw: p.Raw,
})
}
if version == "" {
s.log.DebugContext(ctx, "series unsigned", "capability", "query",
"metric", catalog.ClipMetric(req.Metric))
}
return out, nil
}
// layerNames — тот же словарь `hae.Layers`, переведённый в строки для выборки.
//
// Выводится из словаря, а не перечисляется заново: собственный список разошёлся
// бы с ним молча. Хранилищу перечень нужен по эксплуатационной причине — без
// предиката по слою выборка охватов просматривает все строки метрики за всю
// историю (измерено проходом `ops`: 13.9 мс против 0.026 мс), и цена росла бы
// вместе с возрастом сервиса при любой ширине запроса.
var layerNames = func() []string {
out := make([]string, 0, len(hae.Layers))
for _, l := range hae.Layers {
out = append(out, string(l))
}
return out
}()
// pickLayer выбирает слой по ОХВАТУ точек внутри периода.
//
// Охват — длина пересечения отрезка [первая метка слоя, последняя метка слоя] с
// запрошенным периодом. Слой с пустым пересечением выбывает. Побеждает
// наибольший охват, при равенстве — самый мелкий слой.
//
// Охват меряется метками ТОЧЕК, а не часами объектов, и это не придирка.
// Объекты адресуются часом, поэтому выборка обязана быть шире запроса (точка
// 10:59 живёт в объекте 10:00), а ряд отбирается точной меткой. На периоде
// [10:30, 10:45) часовой слой имеет объект 10:00 с единственной точкой в 10:00,
// минутный — объект 10:00 с точками 10:31…10:44. По часам объектов охваты
// равны, побеждает часовой — и ответ уходит пустым при непустых минутных
// данных. По меткам точек часовой выбывает сразу.
//
// Число точек мерой не является: нижний слой за три плотных дня даёт их больше,
// чем часовой за год, — и «вес за год» вернул бы три дня, не сказав об этом ни
// словом.
func pickLayer(from, to time.Time, spans []store.LayerSpan) string {
best := ""
var bestCover time.Duration
for _, sp := range spans {
if sp.Last.Before(from) || !sp.First.Before(to) {
continue
}
start, end := sp.First, sp.Last
if start.Before(from) {
start = from
}
if end.After(to) {
end = to
}
cover := end.Sub(start)
finer := hae.Rank(hae.Layer(sp.Layer)) < hae.Rank(hae.Layer(best))
if best == "" || cover > bestCover || (cover == bestCover && finer) {
best, bestCover = sp.Layer, cover
}
}
return best
}
// applicable отвечает, можно ли применить объявленный род к ЭТОМУ ряду.
//
// Род — свойство метрики, слой — свойство ряда, и их сочетание бывает опасным.
// Нижний слой HAE это интерполяция, а не сэмплы: сумма по нему завышает втрое
// (находка 34). Конверт, объявляющий `cumulative` рядом с рядом из `raw` и
// молчащий о неприменимости, приглашает потребителя сложить интерполяцию
// самому — система при этом не складывает ничего, а решение у потребителя уже
// принято по завышенному числу.
//
// Слой `sample` под запрет не подпадает: это настоящие сэмплы HealthKit из
// родного экспорта, а не развёртка HAE.
func applicable(style catalog.Style, layer string) bool {
if style == catalog.Unknown || layer == "" {
return false
}
return style != catalog.Cumulative || layer != string(hae.LayerRaw)
}
+132
View File
@@ -0,0 +1,132 @@
package points
import (
"testing"
"time"
"git.vakhrushev.me/av/healthlog/internal/catalog"
"git.vakhrushev.me/av/healthlog/internal/store"
)
func day(d int) time.Time { return time.Date(2026, 6, d, 0, 0, 0, 0, time.UTC) }
func span(layer string, first, last time.Time) store.LayerSpan {
return store.LayerSpan{Layer: layer, First: first, Last: last}
}
// Правило выбора слоя целиком: охват, равенство охватов, выбывание по пустому
// пересечению и вырожденный вход.
func TestВыборСлояПоОхвату(t *testing.T) {
from, to := day(1), day(30)
cases := []struct {
name string
spans []store.LayerSpan
want string
}{
{
// «Вес за год»: нижний слой плотнее, но короче. Победа по числу
// точек вернула бы три дня вместо периода и не сказала бы ни слова.
name: "мелкий слой охватывает меньше крупного",
spans: []store.LayerSpan{
span("raw", day(1), day(4)),
span("hour", day(1), day(29)),
},
want: "hour",
},
{
name: "охваты равны — побеждает мелкий",
spans: []store.LayerSpan{
span("hour", day(1), day(29)),
span("minute", day(1), day(29)),
},
want: "minute",
},
{
// Слой, чьи данные лежат целиком вне периода, выбывает — иначе он
// выиграл бы охват и отдал пустой ряд.
name: "слой лежит вне периода",
spans: []store.LayerSpan{
span("hour", day(1).Add(-48*time.Hour), day(1).Add(-24*time.Hour)),
span("minute", day(2), day(3)),
},
want: "minute",
},
{
// Пересечение считается по ПЕРИОДУ, а не по данным: слой, торчащий
// за обе границы, не получает бесконечного охвата.
name: "слой шире периода",
spans: []store.LayerSpan{
span("hour", day(1).Add(-240*time.Hour), day(30).Add(240*time.Hour)),
span("minute", day(1), day(30)),
},
want: "minute",
},
{
name: "слоёв нет вовсе",
spans: nil,
want: "",
},
{
// Единственная точка ровно на границе периода: охват нулевой, но
// слой пригоден — точка в ответ попадёт.
name: "единственная точка на границе",
spans: []store.LayerSpan{span("raw", from, from)},
want: "raw",
},
{
// Точка ровно на `to` в период не входит (полуинтервал), значит и
// слой из выбора выбывает.
name: "единственная точка на правой границе",
spans: []store.LayerSpan{span("raw", to, to)},
want: "",
},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
if got := pickLayer(from, to, c.spans); got != c.want {
t.Errorf("выбран слой %q, ожидался %q", got, c.want)
}
})
}
}
// Порядок слоёв в выборке не должен влиять на исход: иначе ответ стал бы
// функцией порядка строк, который задаёт SQLite, а не правило.
func TestВыборСлояНеЗависитОтПорядка(t *testing.T) {
from, to := day(1), day(30)
a := span("hour", day(1), day(29))
b := span("minute", day(1), day(29))
if pickLayer(from, to, []store.LayerSpan{a, b}) != pickLayer(from, to, []store.LayerSpan{b, a}) {
t.Error("выбор слоя зависит от порядка охватов")
}
}
// Применимость рода к ОТДАННОМУ ряду. Разряд, ради которого поле заведено:
// накопительная метрика на нижнем слое HAE неприменима, потому что нижний слой
// это интерполяция, а не сэмплы, и сумма по нему завышает втрое.
func TestПрименимостьРодаКРяду(t *testing.T) {
cases := []struct {
style catalog.Style
layer string
want bool
}{
{catalog.Cumulative, "raw", false},
{catalog.Cumulative, "minute", true},
{catalog.Cumulative, "hour", true},
{catalog.Cumulative, "sample", true},
{catalog.Instant, "raw", true},
{catalog.Unknown, "minute", false},
{catalog.Unknown, "raw", false},
{catalog.Cumulative, "", false},
{catalog.Instant, "", false},
}
for _, c := range cases {
if got := applicable(c.style, c.layer); got != c.want {
t.Errorf("применимость %s на слое %q = %v, ожидалось %v", c.style, c.layer, got, c.want)
}
}
}
+300
View File
@@ -0,0 +1,300 @@
package points_test
import (
"context"
"encoding/json"
"log/slog"
"path/filepath"
"sync"
"testing"
"time"
"git.vakhrushev.me/av/healthlog/internal/catalog"
"git.vakhrushev.me/av/healthlog/internal/points"
"git.vakhrushev.me/av/healthlog/internal/store"
)
func service(t *testing.T, in ...store.IncomingPoint) (*points.Service, *store.Store) {
t.Helper()
st, err := store.Open(filepath.Join(t.TempDir(), "healthlog.db"))
if err != nil {
t.Fatalf("store.Open: %v", err)
}
t.Cleanup(func() { _ = st.Close() })
if len(in) > 0 {
if _, err := st.Merge(context.Background(), store.Incoming{Points: in}, store.DeliveryRef{ID: "d"}); err != nil {
t.Fatalf("слияние: %v", err)
}
}
return points.New(st, slog.New(slog.DiscardHandler)), st
}
func hourAt(hh int) time.Time { return time.Date(2026, 6, 1, hh, 0, 0, 0, time.UTC) }
func in(metric, layer string, at time.Time) store.IncomingPoint {
return store.IncomingPoint{
Metric: metric, Layer: layer, Units: "count",
Point: store.Point{Start: at, End: at, Raw: json.RawMessage(`{"qty":1}`)},
}
}
// Версия ответа несёт ГОРИЗОНТ измерения, а не только версию витрины.
//
// Без горизонта метка не меняется, когда час из будущего въезжает в окно сам,
// ходом часов и без единого коммита, — и соседняя задача условного запроса
// подтвердит `304` на ответе, чей род уже перевернулся. Путь проект уже строил
// и закрывал у каталога; здесь он закрывается тем же механизмом.
func TestВерсияОтветаНесётГоризонт(t *testing.T) {
svc, st := service(t, in("m", "raw", hourAt(9)))
ctx := context.Background()
got, err := svc.Series(ctx, points.Request{Metric: "m", From: hourAt(0), To: hourAt(23)})
if err != nil {
t.Fatalf("Series: %v", err)
}
if got.Version == "" {
t.Fatal("ответ без версии — подписывать условный запрос нечем")
}
bare, err := st.StateVersion(ctx)
if err != nil {
t.Fatalf("StateVersion: %v", err)
}
if got.Version == bare {
t.Error("версия ответа равна версии витрины — горизонт в неё не вошёл")
}
if want := catalog.Stamp(bare, catalog.Horizon()); got.Version != want {
t.Errorf("версия ответа %q, ожидалась %q", got.Version, want)
}
}
// Слой выбирается тем же правилом, что проверено на охватах, но уже через
// хранилище: «вес за год» обязан вернуть длинный слой, а не плотный короткий.
func TestРядБерётСлойСНаибольшимОхватом(t *testing.T) {
svc, _ := service(t,
in("body_mass", "raw", hourAt(9)),
in("body_mass", "raw", hourAt(10)),
in("body_mass", "hour", hourAt(1)),
in("body_mass", "hour", hourAt(20)),
)
got, err := svc.Series(context.Background(), points.Request{
Metric: "body_mass", From: hourAt(0), To: hourAt(23),
})
if err != nil {
t.Fatalf("Series: %v", err)
}
if got.Layer != "hour" {
t.Errorf("слой %q, ожидался hour: нижний слой охватывает меньше", got.Layer)
}
if len(got.Points) != 2 {
t.Errorf("точек %d, ожидалось 2", len(got.Points))
}
}
// Явно запрошенный слой уезжает в ответе даже пустым: клиент, спросивший разрез
// поимённо, обязан отличать «за период этого разреза нет» от «параметр
// проигнорирован». Слой, который выбирала система и выбрать не смогла, — пустой.
func TestРядРазличаетПустойЯвныйСлойИОтсутствиеВыбора(t *testing.T) {
svc, _ := service(t, in("m", "raw", hourAt(9)))
ctx := context.Background()
explicit, err := svc.Series(ctx, points.Request{
Metric: "m", From: hourAt(0), To: hourAt(23), Layer: "minute",
})
if err != nil {
t.Fatalf("Series: %v", err)
}
if explicit.Layer != "minute" {
t.Errorf("явный слой %q, ожидался minute", explicit.Layer)
}
if len(explicit.Points) != 0 {
t.Errorf("точек %d, ожидалось 0", len(explicit.Points))
}
chosen, err := svc.Series(ctx, points.Request{
Metric: "нет такой", From: hourAt(0), To: hourAt(23),
})
if err != nil {
t.Fatalf("Series: %v", err)
}
if chosen.Layer != "" {
t.Errorf("слой %q, ожидался пустой: выбирать было не из чего", chosen.Layer)
}
}
// Отмена запроса клиентом — обстоятельство, а не отказ: ответ не собирается, но
// и ERROR владельцу не пишется. Уровень проверяет тест транспорта; здесь —
// что отмена вообще доезжает до драйвера и не игнорируется.
func TestРядУважаетОтменуКонтекста(t *testing.T) {
svc, _ := service(t, in("m", "raw", hourAt(9)))
ctx, cancel := context.WithCancel(context.Background())
cancel()
if _, err := svc.Series(ctx, points.Request{Metric: "m", From: hourAt(0), To: hourAt(23)}); err == nil {
t.Error("отменённый запрос собрал ответ — context до драйвера не доехал")
}
}
// Отказ хранилища доезжает до вызывающего отказом, а не пустым рядом: маршрут
// обязан ответить 500, а не «данных нет». Отказ при этом НЕ транзиентный —
// значит уходит владельцу уровнем ERROR, а не тонет в DEBUG.
func TestРядНаЗакрытомХранилищеОтказывает(t *testing.T) {
svc, st := service(t, in("m", "raw", hourAt(9)))
if err := st.Close(); err != nil {
t.Fatalf("закрытие: %v", err)
}
_, err := svc.Series(context.Background(), points.Request{Metric: "m", From: hourAt(0), To: hourAt(23)})
if err == nil {
t.Fatal("закрытое хранилище отдало ряд")
}
if store.Transient(err) {
t.Error("отказ закрытого хранилища объявлен обстоятельством — владелец о нём не узнает")
}
}
// levels — slog.Handler, копящий уровень и сообщение. Значений атрибутов не
// хранит: проверяется адресат записи, а данные о здоровье в тесты тащить
// незачем.
type levels struct {
mu sync.Mutex
seen []slog.Record
}
func (l *levels) Enabled(context.Context, slog.Level) bool { return true }
func (l *levels) WithAttrs([]slog.Attr) slog.Handler { return l }
func (l *levels) WithGroup(string) slog.Handler { return l }
func (l *levels) Handle(_ context.Context, r slog.Record) error {
l.mu.Lock()
defer l.mu.Unlock()
l.seen = append(l.seen, r.Clone())
return nil
}
func (l *levels) levelOf(msg string) (slog.Level, bool) {
l.mu.Lock()
defer l.mu.Unlock()
for _, r := range l.seen {
if r.Message == msg {
return r.Level, true
}
}
return 0, false
}
func (l *levels) count(level slog.Level) int {
l.mu.Lock()
defer l.mu.Unlock()
n := 0
for _, r := range l.seen {
if r.Level == level {
n++
}
}
return n
}
func loggedService(t *testing.T, in ...store.IncomingPoint) (*points.Service, *store.Store, *levels) {
t.Helper()
st, err := store.Open(filepath.Join(t.TempDir(), "healthlog.db"))
if err != nil {
t.Fatalf("store.Open: %v", err)
}
t.Cleanup(func() { _ = st.Close() })
if len(in) > 0 {
if _, err := st.Merge(context.Background(), store.Incoming{Points: in}, store.DeliveryRef{ID: "d"}); err != nil {
t.Fatalf("слияние: %v", err)
}
}
seen := &levels{}
return points.New(st, slog.New(seen)), st, seen
}
// Отмена клиентом — обстоятельство, а не отказ, и уровень записи это отражает.
//
// Утверждение прямое, потому что иначе оно не держится ничем: смена
// классификации не даёт ни ошибки компиляции, ни красного теста. Агент
// опрашивает маршрут по расписанию, и `ERROR` на каждый его тайм-аут забил бы
// единственный канал, по которому владелец видит настоящий сбой хранилища.
func TestОтменаЗапросаПишетсяDEBUG(t *testing.T) {
svc, _, seen := loggedService(t, in("m", "raw", hourAt(9)))
ctx, cancel := context.WithCancel(context.Background())
cancel()
if _, err := svc.Series(ctx, points.Request{Metric: "m", From: hourAt(0), To: hourAt(23)}); err == nil {
t.Fatal("отменённый запрос собрал ответ")
}
level, ok := seen.levelOf("series interrupted")
if !ok {
t.Fatal("отмена не оставила чекпоинта — исход не наблюдаем")
}
if level != slog.LevelDebug {
t.Errorf("уровень %s, ожидался DEBUG", level)
}
if n := seen.count(slog.LevelError); n != 0 {
t.Errorf("записей ERROR %d, ожидалось 0: отмена клиента — не сбой хранилища", n)
}
}
// Настоящий отказ хранилища доходит до владельца уровнем ERROR.
func TestОтказХранилищаПишетсяERROR(t *testing.T) {
svc, st, seen := loggedService(t, in("m", "raw", hourAt(9)))
if err := st.Close(); err != nil {
t.Fatalf("закрытие: %v", err)
}
if _, err := svc.Series(context.Background(), points.Request{Metric: "m", From: hourAt(0), To: hourAt(23)}); err == nil {
t.Fatal("закрытое хранилище собрало ответ")
}
level, ok := seen.levelOf("series failed")
if !ok {
t.Fatal("отказ не оставил чекпоинта")
}
if level != slog.LevelError {
t.Errorf("уровень %s, ожидался ERROR", level)
}
}
// Предупреждения измерения — привилегия каталога, и маршрут точек их НЕ
// повторяет: агент опрашивает по расписанию, и WARN на каждый опрос обесценил
// бы уровень ровно так же, как обесценила бы его строка на каждый `304`.
func TestМаршрутТочекНеПовторяетПредупрежденияИзмерения(t *testing.T) {
// Метрика с противоречащим родом: часть часов сходится с суммой, часть — со
// средним. Каталог на таком входе пишет WARN.
var seed []store.IncomingPoint
for h := range 8 {
hour := hourAt(h)
coarse := "30"
if h%2 == 0 {
coarse = "15" // среднее двух минутных значений 10 и 20
}
seed = append(seed,
store.IncomingPoint{Metric: "mixed", Layer: "hour", Units: "kJ",
Point: store.Point{Start: hour, End: hour, Raw: json.RawMessage(`{"qty":` + coarse + `}`)}},
store.IncomingPoint{Metric: "mixed", Layer: "minute", Units: "kJ",
Point: store.Point{Start: hour, End: hour, Raw: json.RawMessage(`{"qty":10}`)}},
store.IncomingPoint{Metric: "mixed", Layer: "minute", Units: "kJ",
Point: store.Point{Start: hour.Add(time.Minute), End: hour.Add(time.Minute), Raw: json.RawMessage(`{"qty":20}`)}},
)
}
svc, _, seen := loggedService(t, seed...)
for range 3 {
if _, err := svc.Series(context.Background(), points.Request{
Metric: "mixed", From: hourAt(0), To: hourAt(23),
}); err != nil {
t.Fatalf("Series: %v", err)
}
}
if n := seen.count(slog.LevelWarn); n != 0 {
t.Errorf("маршрут точек написал %d предупреждений — опрос по расписанию обесценит уровень", n)
}
}
+274
View File
@@ -0,0 +1,274 @@
package store
import (
"context"
"database/sql"
"errors"
"fmt"
"sort"
"strings"
"time"
)
// Запросы ряда точек. Вынесены константами по той же причине, что и запросы
// каталога: по ним проверяется план выполнения.
const (
// Объекты выбранного слоя. Точный префикс первичного ключа
// (metric, layer, hour_utc); `payload` здесь и нужен.
seriesPointsQuery = `
SELECT hour_utc, units, payload
FROM bucket
WHERE metric = ? AND layer = ? AND hour_utc BETWEEN ? AND ?
ORDER BY hour_utc`
)
// layerSpansQuery строит выборку охватов слоёв метрики внутри периода.
//
// Границы берутся ТОЧНЫЕ (`first_ts`/`last_ts`), а не по `hour_utc`: объект
// адресуется часом, а ряд отбирается точной меткой, и на периоде короче часа
// эти два множества расходятся. Слой, выбранный по часам, отдал бы пустой ряд
// при непустых данных соседнего слоя — час объекта попадает в период, а его
// единственная точка в период не попадает.
//
// СЛОИ ПЕРЕЧИСЛЕНЫ ЯВНО, и это не украшение запроса, а его цена. Индекс
// `bucket_catalog` идёт `(metric, layer, hour_utc, …)`; без предиката по слою
// SQLite не может сузить поиск по `hour_utc` внутри индекса и просматривает
// ВСЕ строки метрики за всю историю, применяя период построчным фильтром. План
// при этом выглядит успешным (`SEARCH … USING COVERING INDEX`), а цена растёт
// вместе с возрастом сервиса при любой ширине запроса: измерено эксплуатационным
// проходом ревью на копии схемы — 2.06 мс при 52 560 строках метрики против
// 13.9 мс при 350 400, и 0.026 мс с этим перечислением. Словарь слоёв задаёт
// вызывающий: правило принадлежит домену, хранилище лишь выбирает по нему.
func layerSpansQuery(layers int) string {
return `
SELECT layer, min(first_ts), max(last_ts)
FROM bucket
WHERE metric = ? AND layer IN (` +
strings.TrimSuffix(strings.Repeat("?,", layers), ",") + `)
AND hour_utc BETWEEN ? AND ?
GROUP BY layer
ORDER BY layer`
}
// errEmptyWindow — период задан вывернутым. Нарушенный инвариант вызывающего:
// транспорт обязан отвергнуть такой запрос раньше.
var errEmptyWindow = errors.New("окно ряда: from не раньше to")
// errNoLayers — словарь слоёв не задан. Без него выборка охватов выродилась бы
// в скан всей истории метрики, а не отдала бы пустой результат: молчаливая
// деградация цены хуже отказа.
var errNoLayers = errors.New("окно ряда: словарь слоёв пуст")
// LayerSpan — охват одного слоя метрики внутри запрошенного периода.
//
// First и Last — метки первой и последней ТОЧКИ объектов слоя, попавших в
// границы часов запроса. Это границы данных, а не обещание покрытия: внутри
// законно есть дыры.
type LayerSpan struct {
Layer string
First time.Time
Last time.Time
}
// SeriesPoint — точка ряда вместе с единицами объекта, из которого она
// прочитана.
//
// Единицы едут с точкой, а не с рядом: они хранятся на часовом объекте, и
// метрика, чьи объекты разошлись единицами, обязана показать это строкой, а не
// выбрать одно из двух молча.
type SeriesPoint struct {
Point
Units string
}
// SeriesWindow — что читать. Все правила задаёт вызывающий: хранилище выбирает
// строки, а не решает, какие из них правильные.
type SeriesWindow struct {
Metric string
// From включительно, To исключительно.
From, To time.Time
// Layer — явно запрошенный слой; пустая строка означает «выбери сам»,
// и тогда зовётся pick.
Layer string
// Layers — словарь слоёв, среди которых вообще имеет смысл искать. Задаёт
// вызывающий: перечень слоёв — знание домена, а хранилищу он нужен, чтобы
// выборка охватов не превращалась в скан всей истории метрики.
Layers []string
// Measure — окно измерения рода агрегации той же метрики.
Measure CatalogWindow
}
// SeriesSnapshot — весь вход ответа точек, снятый ОДНОЙ транзакцией чтения.
//
// Единый снимок здесь не аккуратность, а условие непротиворечивости: приём идёт
// непрерывно, и фоновая свёртка вправе закоммитить между выбором слоя и чтением
// точек. Тогда слой выбран по одному состоянию витрины, ряд прочитан по
// второму, а род измерен по третьему — ответ внутренне противоречив и от свежего
// неотличим. Пара проб версии такой ответ ОБНАРУЖИТ (метки не будет), но не
// предотвратит: тело всё равно уедет. Тот же довод записан у входа каталога.
type SeriesSnapshot struct {
// Охватов слоёв здесь НЕТ намеренно: они приходят в pick аргументом, и это
// единственные ворота решения о слое. Поле наружу предлагало бы те же данные
// любому будущему вызывающему и приглашало бы решить слой пост-фактум, мимо
// правила, — двое ворот к одному решению.
//
// Layer — слой, из которого собран ряд. Пустая строка означает, что слоя
// нет: выбирать было не из чего либо запрошенный слой пуст.
Layer string
// Points — ряд, отобранный до точных границ периода и упорядоченный.
Points []SeriesPoint
// Pairs — общие часы метрики, от самых свежих к старым.
Pairs []HourPair
}
// ReadSeries снимает вход ответа точек одной транзакцией чтения.
//
// Правило выбора слоя остаётся В ДОМЕНЕ и приходит сюда функцией pick,
// вызываемой ВНУТРИ транзакции. Форма не изобретена: VersionedRead уже
// принимает работу колбэком, а ReadCatalog уже получает параметры правила
// структурой. Перенести само правило сюда значило бы вернуть в хранилище
// решение, которое из него специально убирали.
//
// pick зовётся только когда слой не задан явно и есть из чего выбирать; вернуть
// он может пустую строку — это законный исход «ряда нет».
func (s *Store) ReadSeries(ctx context.Context, w SeriesWindow, pick func([]LayerSpan) string) (SeriesSnapshot, error) {
// Границы В ТЕКСТ ОШИБКИ НЕ ИДУТ. Текст доезжает до записи лога вызывающего,
// а границы периода — часть запроса о здоровье; сегодня транспорт отвергает
// такой запрос раньше, но второй вызывающий (адаптер MCP) откроет этот путь.
if !w.From.Before(w.To) {
return SeriesSnapshot{}, errEmptyWindow
}
if len(w.Layers) == 0 {
return SeriesSnapshot{}, errNoLayers
}
// `ReadOnly` у `modernc.org/sqlite` выбирает `BEGIN` вместо
// `BEGIN IMMEDIATE` и записи НЕ ЗАПРЕЩАЕТ (прочитан исходник драйвера
// зафиксированной версии). То, что этот путь не пишет, держится ревью, а не
// драйвером; полагаться на флаг как на защиту нельзя.
tx, err := s.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true})
if err != nil {
return SeriesSnapshot{}, fmt.Errorf("begin read tx: %w", err)
}
defer func() { _ = tx.Rollback() }()
// Границы по часам ШИРЕ запроса: точка 10:59 живёт в объекте 10:00, и
// огрубление до часа — единственный способ её не потерять. Точный отбор
// идёт ниже, по меткам самих точек.
fromHour := FormatTime(w.From.UTC().Truncate(time.Hour))
toHour := FormatTime(w.To.UTC().Truncate(time.Hour))
spans, err := readLayerSpans(ctx, tx, w, fromHour, toHour)
if err != nil {
return SeriesSnapshot{}, err
}
out := SeriesSnapshot{Layer: w.Layer}
if out.Layer == "" && pick != nil {
out.Layer = pick(spans)
}
if out.Layer != "" {
// Слой ВЫБРАННЫЙ, а не запрошенный: они расходятся ровно тогда, когда
// правило сработало, — то есть в самом частом случае.
if out.Points, err = readSeriesPoints(ctx, tx, w, out.Layer, fromHour, toHour); err != nil {
return SeriesSnapshot{}, err
}
}
pairs, err := commonHours(ctx, tx, w.Metric, w.Measure)
if err != nil {
return SeriesSnapshot{}, err
}
if len(pairs) > 0 {
if err := readHourPairs(ctx, tx, w.Metric, w.Measure, pairs); err != nil {
return SeriesSnapshot{}, err
}
}
out.Pairs = pairs
return out, nil
}
func readLayerSpans(ctx context.Context, tx *sql.Tx, w SeriesWindow, fromHour, toHour string) ([]LayerSpan, error) {
args := make([]any, 0, len(w.Layers)+3)
args = append(args, w.Metric)
for _, l := range w.Layers {
args = append(args, l)
}
args = append(args, fromHour, toHour)
rows, err := tx.QueryContext(ctx, layerSpansQuery(len(w.Layers)), args...)
if err != nil {
return nil, fmt.Errorf("select layer spans: %w", err)
}
defer func() { _ = rows.Close() }()
out := make([]LayerSpan, 0, 4)
for rows.Next() {
var sp LayerSpan
var first, last string
if err := rows.Scan(&sp.Layer, &first, &last); err != nil {
return nil, fmt.Errorf("scan layer span: %w", err)
}
if sp.First, err = ParseTime(first); err != nil {
return nil, err
}
if sp.Last, err = ParseTime(last); err != nil {
return nil, err
}
out = append(out, sp)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("select layer spans: %w", err)
}
return out, nil
}
// readSeriesPoints разжимает объекты выбранного слоя и отбирает точки до точных
// границ периода.
//
// Принадлежность точки периоду определяется её НАЧАЛОМ — тем же правилом, каким
// час объекта берётся по началу точки. Цена названа в спеке: интервал,
// начавшийся раньше from, в ответ не входит.
func readSeriesPoints(ctx context.Context, tx *sql.Tx, w SeriesWindow, layer, fromHour, toHour string) ([]SeriesPoint, error) {
rows, err := tx.QueryContext(ctx, seriesPointsQuery, w.Metric, layer, fromHour, toHour)
if err != nil {
return nil, fmt.Errorf("select series points: %w", err)
}
defer func() { _ = rows.Close() }()
// Непустой срез, а не nil: пустой ряд обязан уехать клиенту как `[]`.
out := make([]SeriesPoint, 0, 64)
for rows.Next() {
var hour, units string
var payload []byte
if err := rows.Scan(&hour, &units, &payload); err != nil {
return nil, fmt.Errorf("scan series bucket: %w", err)
}
points, err := decodePayload(payload)
if err != nil {
return nil, err
}
for _, p := range points {
if p.Start.Before(w.From) || !p.Start.Before(w.To) {
continue
}
out = append(out, SeriesPoint{Point: p, Units: units})
}
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("select series points: %w", err)
}
// Порядок утверждается здесь, а не наследуется от порядка объектов: на нём
// стоят и байтовое утверждение формы ответа, и метка условного запроса.
// Пара (начало, конец) внутри одного слоя одной метрики есть ключ
// идентичности, поэтому порядок ею определён однозначно.
sort.SliceStable(out, func(i, j int) bool {
if !out[i].Start.Equal(out[j].Start) {
return out[i].Start.Before(out[j].Start)
}
return out[i].End.Before(out[j].End)
})
return out, nil
}
+286
View File
@@ -0,0 +1,286 @@
package store_test
import (
"context"
"encoding/json"
"testing"
"time"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// seriesStore заводит витрину с точками, разложенными по слоям.
func seriesStore(t *testing.T, in ...store.IncomingPoint) *store.Store {
t.Helper()
st := open(t)
if _, err := st.Merge(context.Background(), store.Incoming{Points: in}, store.DeliveryRef{ID: "d"}); err != nil {
t.Fatalf("слияние: %v", err)
}
return st
}
// testLayers — словарь слоёв: его задаёт домен, и хранилищу он приходит
// параметром.
var testLayers = []string{"sample", "raw", "minute", "hour", "day"}
func seriesAt(hh, mm int) time.Time {
return time.Date(2026, 6, 1, hh, mm, 0, 0, time.UTC)
}
func seriesPoint(metric, layer string, start time.Time) store.IncomingPoint {
return seriesInterval(metric, layer, start, start)
}
func seriesInterval(metric, layer string, start, end time.Time) store.IncomingPoint {
return store.IncomingPoint{
Metric: metric, Layer: layer, Units: "count",
Point: store.Point{Start: start, End: end, Raw: json.RawMessage(`{"qty":1}`)},
}
}
// readSeries — обёртка с правилом «бери первый попавшийся слой»: правило выбора
// принадлежит домену, а хранилищу проверяется выборка.
func readSeries(t *testing.T, st *store.Store, metric string, from, to time.Time, layer string) store.SeriesSnapshot {
t.Helper()
snap, err := st.ReadSeries(context.Background(), store.SeriesWindow{
Metric: metric, From: from, To: to, Layer: layer, Layers: testLayers,
Measure: store.CatalogWindow{Fine: "minute", Coarse: "hour", Hours: 48, Horizon: seriesAt(23, 0), CoarsePoints: 1, MinFinePoints: 2},
}, func(spans []store.LayerSpan) string {
if len(spans) == 0 {
return ""
}
return spans[0].Layer
})
if err != nil {
t.Fatalf("ReadSeries: %v", err)
}
return snap
}
// Точка на 10:59 живёт в объекте 10:00. Выборка объектов обязана быть ШИРЕ
// запроса, иначе такая точка теряется молча — и потеря видна только тому, кто
// пересчитает точки руками.
func TestРядНеТеряетТочкуВКонцеЧаса(t *testing.T) {
st := seriesStore(t, seriesPoint("m", "minute", seriesAt(10, 59)))
snap := readSeries(t, st, "m", seriesAt(10, 30), seriesAt(11, 0), "minute")
if len(snap.Points) != 1 {
t.Fatalf("точек %d, ожидалась 1", len(snap.Points))
}
if !snap.Points[0].Start.Equal(seriesAt(10, 59)) {
t.Errorf("метка %s, ожидалась 10:59", snap.Points[0].Start)
}
}
// Полуинтервал: точка ровно на `from` в ответе есть, ровно на `to` — нет. Иначе
// два соседних окна посчитали бы граничную точку дважды.
func TestРядБерётГраницыПолуинтервалом(t *testing.T) {
st := seriesStore(t,
seriesPoint("m", "minute", seriesAt(10, 0)),
seriesPoint("m", "minute", seriesAt(11, 0)),
)
snap := readSeries(t, st, "m", seriesAt(10, 0), seriesAt(11, 0), "minute")
if len(snap.Points) != 1 {
t.Fatalf("точек %d, ожидалась 1 (только на from)", len(snap.Points))
}
if !snap.Points[0].Start.Equal(seriesAt(10, 0)) {
t.Errorf("в ответ попала точка %s, а не граница from", snap.Points[0].Start)
}
}
// Охваты меряются метками ТОЧЕК, а не часами объектов, и наблюдается это
// ИСХОДОМ: на периоде короче часа часовой объект попадает в границы часов
// запроса, а его точка в период не попадает. Слой, выбранный по часам, отдал бы
// пустой ряд при непустых минутных данных.
//
// Утверждение стоит на исходе, а не на промежуточных охватах: охваты приходят в
// правило выбора аргументом и наружу не отдаются — двое ворот к одному решению
// однажды разошлись бы.
func TestОхватыСлоёвМеряютсяМеткамиТочек(t *testing.T) {
st := seriesStore(t,
seriesPoint("m", "hour", seriesAt(10, 0)),
seriesPoint("m", "minute", seriesAt(10, 31)),
seriesPoint("m", "minute", seriesAt(10, 44)),
)
var offered []store.LayerSpan
snap, err := st.ReadSeries(context.Background(), store.SeriesWindow{
Metric: "m", From: seriesAt(10, 30), To: seriesAt(10, 45), Layers: testLayers,
Measure: store.CatalogWindow{Fine: "minute", Coarse: "hour", Hours: 48, Horizon: seriesAt(23, 0), CoarsePoints: 1, MinFinePoints: 2},
}, func(spans []store.LayerSpan) string {
offered = spans
// Правило домена в миниатюре: слой, чьи метки лежат вне периода, выбывает.
for _, sp := range spans {
if !sp.Last.Before(seriesAt(10, 30)) && sp.First.Before(seriesAt(10, 45)) {
return sp.Layer
}
}
return ""
})
if err != nil {
t.Fatalf("ReadSeries: %v", err)
}
if len(offered) != 2 {
t.Fatalf("правилу предложено %d слоёв, ожидалось 2", len(offered))
}
if snap.Layer != "minute" {
t.Errorf("выбран слой %q, ожидался minute: у часового нет точек внутри периода", snap.Layer)
}
if len(snap.Points) != 2 {
t.Errorf("точек %d, ожидалось 2 — ряд пуст при непустых данных", len(snap.Points))
}
}
// Пустой период — законный исход, а не отказ. Срез точек при этом непустой:
// nil уехал бы клиенту как `null`.
func TestРядПустогоПериодаОтдаётПустойСрез(t *testing.T) {
st := seriesStore(t, seriesPoint("m", "minute", seriesAt(10, 0)))
snap := readSeries(t, st, "m", seriesAt(20, 0), seriesAt(21, 0), "minute")
if len(snap.Points) != 0 {
t.Fatalf("точек %d, ожидалось 0", len(snap.Points))
}
if snap.Points == nil {
t.Error("срез точек nil — пустой ряд уедет клиенту как null")
}
}
// Ряд упорядочен по (начало, конец) независимо от порядка объектов и точек в
// хранилище: на этом порядке стоят и байтовое утверждение формы, и метка
// условного запроса.
func TestРядУпорядоченПоКоординате(t *testing.T) {
st := seriesStore(t,
seriesPoint("m", "minute", seriesAt(11, 0)),
seriesPoint("m", "minute", seriesAt(10, 0)),
seriesPoint("m", "minute", seriesAt(10, 30)),
)
snap := readSeries(t, st, "m", seriesAt(9, 0), seriesAt(12, 0), "minute")
if len(snap.Points) != 3 {
t.Fatalf("точек %d, ожидалось 3", len(snap.Points))
}
for i := 1; i < len(snap.Points); i++ {
if snap.Points[i].Start.Before(snap.Points[i-1].Start) {
t.Fatalf("порядок нарушен: %s после %s", snap.Points[i].Start, snap.Points[i-1].Start)
}
}
}
// Вывернутый период — отказ хранилища, а не пустой ответ: такой запрос до
// витрины доходить не должен, и если дошёл, молчать об этом нельзя.
func TestРядОтвергаетВывернутыйПериод(t *testing.T) {
st := seriesStore(t, seriesPoint("m", "minute", seriesAt(10, 0)))
_, err := st.ReadSeries(context.Background(), store.SeriesWindow{
Metric: "m", From: seriesAt(11, 0), To: seriesAt(10, 0), Layers: testLayers,
Measure: store.CatalogWindow{Fine: "minute", Coarse: "hour", Hours: 48, Horizon: seriesAt(23, 0)},
}, nil)
if err == nil {
t.Error("вывернутый период принят молча")
}
}
// Закрытое хранилище — отказ, а не пустой ряд. Ветка редкая, но молчащая:
// маршрут чтения, получивший пустой ряд вместо ошибки, отдал бы клиенту
// «данных нет» на остановленном сервисе.
func TestРядНаЗакрытомХранилищеОтказывает(t *testing.T) {
st := seriesStore(t, seriesPoint("m", "minute", seriesAt(10, 0)))
if err := st.Close(); err != nil {
t.Fatalf("закрытие: %v", err)
}
_, err := st.ReadSeries(context.Background(), store.SeriesWindow{
Metric: "m", From: seriesAt(9, 0), To: seriesAt(11, 0), Layer: "minute", Layers: testLayers,
Measure: store.CatalogWindow{Fine: "minute", Coarse: "hour", Hours: 48, Horizon: seriesAt(23, 0)},
}, nil)
if err == nil {
t.Error("закрытое хранилище отдало ряд")
}
}
// Порядок при СОВПАДАЮЩЕМ начале решается концом координаты. Ветка достижима
// только на интервальных точках — а под одной меткой лежит до трёх записей сна,
// и на их порядке стоят и байтовое утверждение формы ответа, и метка условного
// запроса.
func TestРядУпорядоченПоКонцуПриРавномНачале(t *testing.T) {
start := seriesAt(1, 0)
st := seriesStore(t,
seriesInterval("sleep_analysis", "raw", start, seriesAt(3, 30)),
seriesInterval("sleep_analysis", "raw", start, seriesAt(2, 0)),
seriesInterval("sleep_analysis", "raw", start, seriesAt(2, 45)),
)
snap := readSeries(t, st, "sleep_analysis", seriesAt(0, 0), seriesAt(5, 0), "raw")
if len(snap.Points) != 3 {
t.Fatalf("точек %d, ожидалось 3: одна метка, разные интервалы", len(snap.Points))
}
for i, want := range []time.Time{seriesAt(2, 0), seriesAt(2, 45), seriesAt(3, 30)} {
if !snap.Points[i].Start.Equal(start) {
t.Fatalf("точка %d начинается в %s, а не в общей метке", i, snap.Points[i].Start)
}
if !snap.Points[i].End.Equal(want) {
t.Errorf("точка %d кончается в %s, ожидалось %s", i, snap.Points[i].End, want)
}
}
}
// ОДИН СНИМОК ВИТРИНЫ на весь ответ, а не три чтения подряд.
//
// Приём идёт непрерывно, и фоновая свёртка вправе закоммитить между выбором
// слоя и чтением точек. Тогда слой выбран по одному состоянию, ряд прочитан по
// второму, а род измерен по третьему — ответ внутренне противоречив и от
// свежего неотличим. Пара проб версии такой ответ обнаружит (метки не будет),
// но НЕ предотвратит: тело всё равно уедет.
//
// Оракул точный: колбэк выбора слоя — готовая точка вклинивания, и коммит из
// него идёт другим соединением, пока транзакция чтения открыта.
func TestРядСобранИзОдногоСнимкаВитрины(t *testing.T) {
st := seriesStore(t, seriesPoint("m", "minute", seriesAt(10, 0)))
ctx := context.Background()
committed := false
snap, err := st.ReadSeries(ctx, store.SeriesWindow{
Metric: "m", From: seriesAt(9, 0), To: seriesAt(12, 0), Layers: testLayers,
Measure: store.CatalogWindow{Fine: "minute", Coarse: "hour", Hours: 48, Horizon: seriesAt(23, 0), CoarsePoints: 1, MinFinePoints: 2},
}, func(spans []store.LayerSpan) string {
// Свёртка коммитит ВНУТРИ чтения — ровно то, что делает воркер.
if _, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{
seriesPoint("m", "minute", seriesAt(10, 30)),
}}, store.DeliveryRef{ID: "late"}); err != nil {
t.Errorf("слияние во время чтения: %v", err)
return ""
}
committed = true
return "minute"
})
if err != nil {
t.Fatalf("ReadSeries: %v", err)
}
if !committed {
t.Fatal("коммит во время чтения не состоялся — утверждение ничего не доказывает")
}
if len(snap.Points) != 1 {
t.Fatalf("точек %d, ожидалась 1: ряд собран из двух состояний витрины", len(snap.Points))
}
if !snap.Points[0].Start.Equal(seriesAt(10, 0)) {
t.Errorf("в ряд попала точка %s, закоммиченная уже после выбора слоя", snap.Points[0].Start)
}
}
// Пустой словарь слоёв — отказ, а не пустой ответ. Без предиката по слою
// выборка охватов просматривает все строки метрики за всю историю, и молчаливая
// деградация цены хуже отказа: план запроса при этом выглядит успешным.
func TestРядТребуетСловарьСлоёв(t *testing.T) {
st := seriesStore(t, seriesPoint("m", "minute", seriesAt(10, 0)))
_, err := st.ReadSeries(context.Background(), store.SeriesWindow{
Metric: "m", From: seriesAt(9, 0), To: seriesAt(11, 0),
Measure: store.CatalogWindow{Fine: "minute", Coarse: "hour", Hours: 48, Horizon: seriesAt(23, 0)},
}, nil)
if err == nil {
t.Error("пустой словарь слоёв принят молча")
}
}