добавлен приём пакетов Health Auto Export

- POST /api/v1/ingest: токен, лимит тела, gzip, durable-запись тела в сырой
  архив, учёт доставки в SQLite
- код ответа отражает доставку, а не разбор: битый JSON — 400, непонятое
  содержимое — 200, данные уже сохранены и доразберутся позже
- сохраняется полный набор заголовков запроса с вычисткой секретов по имени
  и по совпадению значения с токеном
This commit is contained in:
av
2026-08-01 12:37:17 +03:00
parent 5e2385ba6e
commit 41f2f05f2f
20 changed files with 1903 additions and 0 deletions
+55
View File
@@ -0,0 +1,55 @@
package main
import (
"context"
"flag"
"fmt"
"net/http"
"strings"
"time"
"git.vakhrushev.me/av/healthlog/internal/config"
)
func runHealthcheck(args []string) error {
fs := flag.NewFlagSet("healthcheck", flag.ContinueOnError)
cfgPath := fs.String("config", config.DefaultPath, "путь к config.toml")
timeout := fs.Duration("timeout", 5*time.Second, "таймаут запроса")
if err := fs.Parse(args); err != nil {
return fmt.Errorf("parse flags: %w", err)
}
cfg, err := config.Load(*cfgPath)
if err != nil {
return err
}
ctx, cancel := context.WithTimeout(context.Background(), *timeout)
defer cancel()
url := "http://" + localAddr(cfg.Server.Addr) + "/healthz"
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return fmt.Errorf("build request: %w", err)
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
return fmt.Errorf("healthz: %w", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("healthz: статус %d", resp.StatusCode)
}
return nil
}
// localAddr превращает адрес прослушивания в адрес для обращения к себе:
// ":8080" слушает все интерфейсы, но стучаться надо в конкретный.
func localAddr(addr string) string {
if strings.HasPrefix(addr, ":") {
return "127.0.0.1" + addr
}
return addr
}
+48
View File
@@ -0,0 +1,48 @@
// Команда healthlog — коллектор данных Apple Health из Health Auto Export.
//
// Подкоманды:
//
// healthlog [serve] --config <path> принимать пакеты (по умолчанию)
// healthlog healthcheck --config <p> проверить /healthz (для docker HEALTHCHECK)
package main
import (
"os"
"strings"
"git.vakhrushev.me/av/healthlog/internal/logging"
)
func main() {
args := os.Args[1:]
// Первый позиционный аргумент (не флаг) — подкоманда. Без него (и при
// `--config ...`) запускаем сервис.
cmd := "serve"
if len(args) > 0 && !strings.HasPrefix(args[0], "-") {
cmd, args = args[0], args[1:]
}
var err error
switch cmd {
case "serve":
err = runServe(args)
case "healthcheck":
err = runHealthcheck(args)
default:
_, _ = os.Stderr.WriteString("unknown command: " + cmd + "\n")
os.Exit(2)
}
if err != nil {
if cmd == "serve" {
// Фатальный сбой старта сервиса — структурный лог ERROR (как в
// проде), затем ненулевой код возврата.
logging.NewStderr().Error("fatal startup", "command", cmd, "error", err)
} else {
// Диагностический CLI — человекочитаемый stderr, это
// пользовательский вывод, а не лог сервиса.
_, _ = os.Stderr.WriteString("fatal: " + err.Error() + "\n")
}
os.Exit(1)
}
}
+99
View File
@@ -0,0 +1,99 @@
package main
import (
"context"
"errors"
"flag"
"fmt"
"net/http"
"os/signal"
"syscall"
"time"
"git.vakhrushev.me/av/healthlog/internal/archive"
"git.vakhrushev.me/av/healthlog/internal/config"
"git.vakhrushev.me/av/healthlog/internal/httpapi"
"git.vakhrushev.me/av/healthlog/internal/ingest"
"git.vakhrushev.me/av/healthlog/internal/logging"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// shutdownTimeout — сколько ждём завершения активных запросов при остановке.
// Приём может быть в середине записи многомегабайтного тела в архив.
const shutdownTimeout = 30 * time.Second
func runServe(args []string) error {
fs := flag.NewFlagSet("serve", flag.ContinueOnError)
cfgPath := fs.String("config", config.DefaultPath, "путь к config.toml")
if err := fs.Parse(args); err != nil {
return fmt.Errorf("parse flags: %w", err)
}
cfg, err := config.Load(*cfgPath)
if err != nil {
return err
}
log := logging.New(cfg.Log.Level, cfg.Log.Format)
st, err := store.Open(cfg.Storage.DBPath)
if err != nil {
return err
}
defer func() { _ = st.Close() }()
arch, err := archive.New(cfg.Storage.ArchiveDir)
if err != nil {
return err
}
if len(cfg.Auth.WriteTokens) == 0 {
log.Warn("write auth disabled", "reason", "auth.write_tokens пуст")
}
handler := httpapi.New(httpapi.Options{
Ingest: ingest.New(arch, st, log),
Log: log,
WriteTokens: cfg.Auth.WriteTokens,
MaxBodyMB: cfg.Ingest.MaxBodyMB,
})
srv := &http.Server{
Addr: cfg.Server.Addr,
Handler: handler,
// ReadTimeout щедрый (большой пакет по мобильной сети), но заголовки
// обязаны приехать быстро — иначе полуоткрытое соединение держит слот.
ReadHeaderTimeout: 10 * time.Second,
ReadTimeout: cfg.Server.ReadTimeout.D(),
WriteTimeout: cfg.Server.WriteTimeout.D(),
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
errCh := make(chan error, 1)
go func() {
log.Info("server started",
"addr", cfg.Server.Addr,
"db_path", cfg.Storage.DBPath,
"archive_dir", arch.Root(),
"max_body_mb", cfg.Ingest.MaxBodyMB)
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
errCh <- fmt.Errorf("listen: %w", err)
}
}()
select {
case err := <-errCh:
return err
case <-ctx.Done():
}
log.Info("server stopping")
shutdownCtx, cancel := context.WithTimeout(context.Background(), shutdownTimeout)
defer cancel()
if err := srv.Shutdown(shutdownCtx); err != nil {
return fmt.Errorf("shutdown: %w", err)
}
return nil
}
+129
View File
@@ -0,0 +1,129 @@
// Package archive — сырой архив принятых тел запросов.
//
// Архив, а не витрина, является источником истины: тело каждого запроса
// ложится на диск как есть до любого разбора, и витрина в любой момент
// пересобирается из него. Отсюда требование к записи — durability: пока файл
// не лёг на диск, доставка не считается принятой.
package archive
import (
"compress/gzip"
"fmt"
"io"
"os"
"path/filepath"
"time"
)
// Archive — каталог сырых тел, разложенных по дате приёма.
type Archive struct {
root string
}
// New создаёт архив в указанном каталоге (каталог создаётся при отсутствии).
func New(root string) (*Archive, error) {
if err := os.MkdirAll(root, 0o755); err != nil {
return nil, fmt.Errorf("create archive dir %q: %w", root, err)
}
return &Archive{root: root}, nil
}
// Root возвращает корневой каталог архива.
func (a *Archive) Root() string { return a.root }
// Write сохраняет тело запроса и возвращает путь относительно корня архива.
//
// Запись durable: временный файл → fsync → атомарный rename → fsync каталога.
// Без этого «сохранили — значит приняли» было бы обещанием, которое не
// переживает потерю питания.
func (a *Archive) Write(id string, at time.Time, body []byte) (string, error) {
rel := filepath.Join(at.UTC().Format("2006/01/02"), id+".json.gz")
full := filepath.Join(a.root, rel)
dir := filepath.Dir(full)
if err := os.MkdirAll(dir, 0o755); err != nil {
return "", fmt.Errorf("create archive dir %q: %w", dir, err)
}
tmp := full + ".tmp"
if err := writeGzip(tmp, body); err != nil {
_ = os.Remove(tmp)
return "", err
}
if err := os.Rename(tmp, full); err != nil {
_ = os.Remove(tmp)
return "", fmt.Errorf("rename archive file: %w", err)
}
if err := syncDir(dir); err != nil {
return "", err
}
return rel, nil
}
// Open открывает сохранённое тело для чтения (распакованным). Нужен для
// пересборки витрины из архива.
func (a *Archive) Open(rel string) (io.ReadCloser, error) {
f, err := os.Open(filepath.Join(a.root, rel))
if err != nil {
return nil, fmt.Errorf("open archive file %q: %w", rel, err)
}
gz, err := gzip.NewReader(f)
if err != nil {
_ = f.Close()
return nil, fmt.Errorf("gunzip archive file %q: %w", rel, err)
}
return readCloser{gz: gz, f: f}, nil
}
type readCloser struct {
gz *gzip.Reader
f *os.File
}
func (r readCloser) Read(p []byte) (int, error) { return r.gz.Read(p) } //nolint:wrapcheck // io.Reader: ошибку отдаём как есть
func (r readCloser) Close() error {
err := r.gz.Close()
if cerr := r.f.Close(); err == nil {
err = cerr
}
if err != nil {
return fmt.Errorf("close archive file: %w", err)
}
return nil
}
func writeGzip(path string, body []byte) error {
f, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0o644)
if err != nil {
return fmt.Errorf("create archive file: %w", err)
}
defer func() { _ = f.Close() }()
gz := gzip.NewWriter(f)
if _, err := gz.Write(body); err != nil {
return fmt.Errorf("gzip archive body: %w", err)
}
if err := gz.Close(); err != nil {
return fmt.Errorf("flush gzip: %w", err)
}
if err := f.Sync(); err != nil {
return fmt.Errorf("fsync archive file: %w", err)
}
return nil
}
// syncDir фиксирует появление нового имени в каталоге: без этого rename может
// не пережить внезапную перезагрузку.
func syncDir(dir string) error {
d, err := os.Open(dir)
if err != nil {
return fmt.Errorf("open archive dir: %w", err)
}
defer func() { _ = d.Close() }()
if err := d.Sync(); err != nil {
return fmt.Errorf("fsync archive dir: %w", err)
}
return nil
}
+92
View File
@@ -0,0 +1,92 @@
package archive_test
import (
"io"
"os"
"path/filepath"
"strings"
"testing"
"time"
"git.vakhrushev.me/av/healthlog/internal/archive"
)
func TestWriteAndOpenRoundTrip(t *testing.T) {
a := newArchive(t)
body := []byte(`{"data":{"metrics":[]}}`)
at := time.Date(2026, 7, 31, 14, 5, 0, 0, time.UTC)
rel, err := a.Write("01k1abc", at, body)
if err != nil {
t.Fatalf("Write: %v", err)
}
if want := filepath.Join("2026", "07", "31", "01k1abc.json.gz"); rel != want {
t.Errorf("путь = %q, ожидался %q", rel, want)
}
r, err := a.Open(rel)
if err != nil {
t.Fatalf("Open: %v", err)
}
defer func() { _ = r.Close() }()
got, err := io.ReadAll(r)
if err != nil {
t.Fatalf("чтение: %v", err)
}
if string(got) != string(body) {
t.Errorf("прочитано %q, записано %q", got, body)
}
}
// Раскладка по дате приёма — в UTC: иначе один и тот же момент попадал бы в
// разные каталоги в зависимости от зоны машины.
func TestWriteUsesUTCDate(t *testing.T) {
a := newArchive(t)
// 01:30 по московскому времени — это ещё предыдущие сутки по UTC.
msk := time.FixedZone("MSK", 3*60*60)
at := time.Date(2026, 8, 1, 1, 30, 0, 0, msk)
rel, err := a.Write("01k1def", at, []byte(`{"data":{}}`))
if err != nil {
t.Fatalf("Write: %v", err)
}
if want := filepath.Join("2026", "07", "31", "01k1def.json.gz"); rel != want {
t.Errorf("путь = %q, ожидался %q", rel, want)
}
}
// Временных файлов после успешной записи остаться не должно.
func TestWriteLeavesNoTempFiles(t *testing.T) {
root := t.TempDir()
a, err := archive.New(root)
if err != nil {
t.Fatalf("New: %v", err)
}
if _, err := a.Write("01k1ghi", time.Date(2026, 7, 31, 0, 0, 0, 0, time.UTC), []byte(`{"data":{}}`)); err != nil {
t.Fatalf("Write: %v", err)
}
err = filepath.Walk(root, func(path string, info os.FileInfo, err error) error {
if err != nil {
return err
}
if strings.HasSuffix(path, ".tmp") {
t.Errorf("остался временный файл: %s", path)
}
return nil
})
if err != nil {
t.Fatalf("обход архива: %v", err)
}
}
func newArchive(t *testing.T) *archive.Archive {
t.Helper()
a, err := archive.New(filepath.Join(t.TempDir(), "raw"))
if err != nil {
t.Fatalf("New: %v", err)
}
return a
}
+176
View File
@@ -0,0 +1,176 @@
// Package config загружает конфигурацию healthlog из TOML-файла.
package config
import (
"errors"
"fmt"
"os"
"strings"
"time"
"github.com/pelletier/go-toml/v2"
)
// DefaultPath — имя конфига по умолчанию: ищется в рабочей директории
// процесса. Переопределяется опцией --config=path.
const DefaultPath = "config.toml"
// Config — корневая конфигурация сервиса (см. config.example.toml).
type Config struct {
Server Server `toml:"server"`
Auth Auth `toml:"auth"`
Storage Storage `toml:"storage"`
Ingest Ingest `toml:"ingest"`
Log Log `toml:"log"`
}
// Server — параметры HTTP-сервера.
type Server struct {
Addr string `toml:"addr"`
// ReadTimeout — на всё чтение запроса вместе с телом. Держим щедрым:
// экспорт истории — это десятки мегабайт по мобильной сети.
ReadTimeout Duration `toml:"read_timeout"`
WriteTimeout Duration `toml:"write_timeout"`
}
// Auth — токены доступа. Пустой список = проверка выключена (локальный
// запуск в доверенной сети); сервис предупреждает об этом на старте.
type Auth struct {
WriteTokens []string `toml:"write_tokens"`
ReadTokens []string `toml:"read_tokens"`
}
// Storage — где лежат сырой архив и витрина.
type Storage struct {
DBPath string `toml:"db_path"`
ArchiveDir string `toml:"archive_dir"`
}
// Ingest — параметры приёма.
type Ingest struct {
MaxBodyMB int `toml:"max_body_mb"`
}
// Log — уровень и формат логов.
type Log struct {
Level string `toml:"level"`
Format string `toml:"format"`
}
// Load читает TOML-конфиг, подставляет умолчания и валидирует результат.
// Отсутствующий файл — не ошибка: сервис поднимается на умолчаниях, что
// нужно для локального запуска.
func Load(path string) (*Config, error) {
var cfg Config
data, err := os.ReadFile(path)
switch {
case errors.Is(err, os.ErrNotExist):
// Конфига нет — работаем на умолчаниях.
case err != nil:
return nil, fmt.Errorf("read config: %w", err)
default:
if err := toml.Unmarshal(data, &cfg); err != nil {
return nil, fmt.Errorf("parse config: %w", err)
}
}
cfg.applyDefaults()
if err := cfg.validate(); err != nil {
return nil, err
}
return &cfg, nil
}
func (c *Config) applyDefaults() {
if c.Server.Addr == "" {
c.Server.Addr = ":8080"
}
if c.Server.ReadTimeout == 0 {
c.Server.ReadTimeout = Duration(5 * time.Minute)
}
if c.Server.WriteTimeout == 0 {
c.Server.WriteTimeout = Duration(30 * time.Second)
}
if c.Storage.DBPath == "" {
c.Storage.DBPath = "./healthlog.db"
}
if c.Storage.ArchiveDir == "" {
c.Storage.ArchiveDir = "./raw"
}
if c.Ingest.MaxBodyMB == 0 {
c.Ingest.MaxBodyMB = 64
}
if c.Log.Level == "" {
c.Log.Level = "info"
}
if c.Log.Format == "" {
c.Log.Format = "json"
}
}
// validate собирает все проблемы разом (errors.Join), чтобы не чинить конфиг
// по одной строчке за перезапуск.
func (c *Config) validate() error {
var errs []error
if c.Ingest.MaxBodyMB <= 0 {
errs = append(errs, fmt.Errorf("ingest.max_body_mb: %w: %d", ErrOutOfRange, c.Ingest.MaxBodyMB))
}
if c.Server.ReadTimeout <= 0 {
errs = append(errs, fmt.Errorf("server.read_timeout: %w", ErrOutOfRange))
}
if c.Server.WriteTimeout <= 0 {
errs = append(errs, fmt.Errorf("server.write_timeout: %w", ErrOutOfRange))
}
if !oneOf(c.Log.Level, "debug", "info", "warn", "error") {
errs = append(errs, fmt.Errorf("log.level: %w: %q", ErrUnknownValue, c.Log.Level))
}
if !oneOf(c.Log.Format, "json", "text") {
errs = append(errs, fmt.Errorf("log.format: %w: %q", ErrUnknownValue, c.Log.Format))
}
for i, t := range c.Auth.WriteTokens {
if strings.TrimSpace(t) == "" {
errs = append(errs, fmt.Errorf("auth.write_tokens[%d]: %w", i, ErrEmptyToken))
}
}
for i, t := range c.Auth.ReadTokens {
if strings.TrimSpace(t) == "" {
errs = append(errs, fmt.Errorf("auth.read_tokens[%d]: %w", i, ErrEmptyToken))
}
}
return errors.Join(errs...)
}
// Ошибки валидации конфига.
var (
ErrOutOfRange = errors.New("значение вне допустимого диапазона")
ErrUnknownValue = errors.New("недопустимое значение")
ErrEmptyToken = errors.New("пустой токен")
)
func oneOf(v string, allowed ...string) bool {
for _, a := range allowed {
if strings.EqualFold(v, a) {
return true
}
}
return false
}
// Duration — time.Duration, читаемая из TOML строкой Go-формата ("30s", "5m").
type Duration time.Duration
// UnmarshalText разбирает Go-длительность из TOML-строки.
func (d *Duration) UnmarshalText(text []byte) error {
v, err := time.ParseDuration(string(text))
if err != nil {
return fmt.Errorf("parse duration %q: %w", text, err)
}
*d = Duration(v)
return nil
}
// D возвращает значение как time.Duration.
func (d Duration) D() time.Duration { return time.Duration(d) }
+111
View File
@@ -0,0 +1,111 @@
package config_test
import (
"errors"
"os"
"path/filepath"
"testing"
"time"
"git.vakhrushev.me/av/healthlog/internal/config"
)
// Отсутствующий конфиг — не ошибка: локальный запуск должен работать «из
// коробки», без файла.
func TestLoadMissingFileUsesDefaults(t *testing.T) {
cfg, err := config.Load(filepath.Join(t.TempDir(), "нет-такого.toml"))
if err != nil {
t.Fatalf("Load: %v", err)
}
if cfg.Server.Addr != ":8080" {
t.Errorf("addr = %q, ожидалось \":8080\"", cfg.Server.Addr)
}
if cfg.Ingest.MaxBodyMB != 64 {
t.Errorf("max_body_mb = %d, ожидалось 64", cfg.Ingest.MaxBodyMB)
}
if cfg.Server.ReadTimeout.D() != 5*time.Minute {
t.Errorf("read_timeout = %v, ожидалось 5m", cfg.Server.ReadTimeout.D())
}
if len(cfg.Auth.WriteTokens) != 0 {
t.Errorf("write_tokens = %v, ожидался пустой список", cfg.Auth.WriteTokens)
}
}
func TestLoadReadsFile(t *testing.T) {
path := writeConfig(t, `
[server]
addr = "127.0.0.1:9000"
read_timeout = "1m"
[auth]
write_tokens = ["секрет"]
[ingest]
max_body_mb = 8
[log]
level = "debug"
format = "text"
`)
cfg, err := config.Load(path)
if err != nil {
t.Fatalf("Load: %v", err)
}
if cfg.Server.Addr != "127.0.0.1:9000" {
t.Errorf("addr = %q", cfg.Server.Addr)
}
if cfg.Server.ReadTimeout.D() != time.Minute {
t.Errorf("read_timeout = %v", cfg.Server.ReadTimeout.D())
}
if cfg.Ingest.MaxBodyMB != 8 {
t.Errorf("max_body_mb = %d", cfg.Ingest.MaxBodyMB)
}
if len(cfg.Auth.WriteTokens) != 1 || cfg.Auth.WriteTokens[0] != "секрет" {
t.Errorf("write_tokens = %v", cfg.Auth.WriteTokens)
}
}
// Валидация собирает все проблемы разом (errors.Join), а не чинится по одной
// за перезапуск.
func TestValidateReportsAllProblems(t *testing.T) {
path := writeConfig(t, `
[ingest]
max_body_mb = -1
[auth]
write_tokens = [" "]
[log]
level = "нет-такого"
format = "yaml"
`)
_, err := config.Load(path)
if err == nil {
t.Fatal("ожидалась ошибка валидации")
}
for _, want := range []error{config.ErrOutOfRange, config.ErrUnknownValue, config.ErrEmptyToken} {
if !errors.Is(err, want) {
t.Errorf("ошибка не содержит %v: %v", want, err)
}
}
}
func TestLoadRejectsBrokenTOML(t *testing.T) {
if _, err := config.Load(writeConfig(t, "это не toml =")); err == nil {
t.Fatal("ожидалась ошибка разбора")
}
}
func writeConfig(t *testing.T, body string) string {
t.Helper()
path := filepath.Join(t.TempDir(), "config.toml")
if err := os.WriteFile(path, []byte(body), 0o600); err != nil {
t.Fatalf("подготовка конфига: %v", err)
}
return path
}
+138
View File
@@ -0,0 +1,138 @@
// Package httpapi — HTTP-транспорт healthlog: приём пакетов и (позже) read API.
//
// Транспорт тонкий: разбирает запрос, зовёт use-case, переводит его ошибку в
// ответ. Исход операции логирует use-case, а не транспорт.
package httpapi
import (
"crypto/subtle"
"encoding/json"
"log/slog"
"net/http"
"time"
"github.com/go-chi/chi/v5"
"github.com/go-chi/chi/v5/middleware"
"git.vakhrushev.me/av/healthlog/internal/ingest"
)
// Options — зависимости и настройки транспорта.
type Options struct {
Ingest *ingest.Service
Log *slog.Logger
WriteTokens []string
MaxBodyMB int
}
type api struct {
ingest *ingest.Service
log *slog.Logger
writeTokens []string
maxBody int64
}
// New собирает HTTP-роутер.
func New(o Options) http.Handler {
a := &api{
ingest: o.Ingest,
log: o.Log,
writeTokens: o.WriteTokens,
maxBody: int64(o.MaxBodyMB) << 20,
}
r := chi.NewRouter()
r.Use(middleware.Recoverer)
r.Use(a.accessLog)
r.Get("/healthz", a.handleHealthz)
r.Route("/api/v1", func(r chi.Router) {
r.With(a.requireWriteToken).Post("/ingest", a.handleIngest)
})
return r
}
func (a *api) handleHealthz(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"status":"ok"}`))
}
// requireWriteToken проверяет токен приёма. Пустой список токенов = проверка
// выключена: локальный запуск в доверенной сети. О выключенной проверке
// сервис предупреждает на старте.
func (a *api) requireWriteToken(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if len(a.writeTokens) == 0 {
next.ServeHTTP(w, r)
return
}
if !tokenAllowed(bearer(r), a.writeTokens) {
writeError(w, http.StatusUnauthorized, "неверный или отсутствующий токен")
return
}
next.ServeHTTP(w, r)
})
}
func bearer(r *http.Request) string {
const prefix = "Bearer "
h := r.Header.Get("Authorization")
if len(h) > len(prefix) && h[:len(prefix)] == prefix {
return h[len(prefix):]
}
return ""
}
// tokenAllowed сравнивает токен за постоянное время: побайтовое сравнение с
// ранним выходом утекает длину совпадающего префикса.
func tokenAllowed(got string, allowed []string) bool {
ok := false
for _, want := range allowed {
if subtle.ConstantTimeCompare([]byte(got), []byte(want)) == 1 {
ok = true
}
}
return ok
}
// accessLog пишет одну запись на запрос. Рутинно-частые эндпоинты
// (healthcheck) — на DEBUG, чтобы не забивать аудит.
func (a *api) accessLog(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
start := time.Now()
ww := middleware.NewWrapResponseWriter(w, r.ProtoMajor)
next.ServeHTTP(ww, r)
level := slog.LevelInfo
if r.URL.Path == "/healthz" {
level = slog.LevelDebug
}
a.log.Log(r.Context(), level, "http request",
"transport", "http",
"http.method", r.Method,
"http.route", routePattern(r),
"http.status_code", ww.Status(),
"duration_ms", time.Since(start).Milliseconds())
})
}
func routePattern(r *http.Request) string {
if rctx := chi.RouteContext(r.Context()); rctx != nil {
if p := rctx.RoutePattern(); p != "" {
return p
}
}
return r.URL.Path
}
func writeJSON(w http.ResponseWriter, status int, v any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(v)
}
func writeError(w http.ResponseWriter, status int, msg string) {
writeJSON(w, status, map[string]string{"error": msg})
}
+244
View File
@@ -0,0 +1,244 @@
package httpapi_test
import (
"bytes"
"compress/gzip"
"encoding/json"
"log/slog"
"net/http"
"net/http/httptest"
"path/filepath"
"strings"
"testing"
"git.vakhrushev.me/av/healthlog/internal/archive"
"git.vakhrushev.me/av/healthlog/internal/httpapi"
"git.vakhrushev.me/av/healthlog/internal/ingest"
"git.vakhrushev.me/av/healthlog/internal/store"
)
const samplePayload = `{"data":{"metrics":[{"name":"step_count","units":"count","data":[{"qty":8,"date":"2026-07-31 12:00:00 +0300"}]}],"workouts":[]}}`
func TestIngestAcceptsPayload(t *testing.T) {
h, st := newAPI(t, nil)
req := httptest.NewRequest(http.MethodPost, "/api/v1/ingest", strings.NewReader(samplePayload))
req.Header.Set("Content-Type", "application/json")
req.Header.Set("automation-name", "healthlog")
req.Header.Set("automation-aggregation", "None")
req.Header.Set("session-id", "sess-1")
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("статус = %d, тело = %s", rec.Code, rec.Body.String())
}
var resp struct {
DeliveryID string `json:"delivery_id"`
Bytes int64 `json:"bytes"`
SHA256 string `json:"sha256"`
}
if err := json.Unmarshal(rec.Body.Bytes(), &resp); err != nil {
t.Fatalf("разбор ответа: %v", err)
}
if resp.DeliveryID == "" || resp.SHA256 == "" {
t.Errorf("ответ без идентификатора или хеша: %+v", resp)
}
if resp.Bytes != int64(len(samplePayload)) {
t.Errorf("bytes = %d, ожидалось %d", resp.Bytes, len(samplePayload))
}
if n, err := st.CountDeliveries(t.Context()); err != nil || n != 1 {
t.Errorf("доставок = %d (err=%v), ожидалась 1", n, err)
}
}
// Заголовки сохраняются целиком: документация HAE неполна, и незадокументированные
// поля уже приносили самые полезные находки.
func TestIngestStoresAllHeaders(t *testing.T) {
h, st := newAPI(t, nil)
req := httptest.NewRequest(http.MethodPost, "/api/v1/ingest", strings.NewReader(samplePayload))
req.Header.Set("automation-id", "BC99C8A3-8BE7-4519-B545-F3ED6212008E")
req.Header.Set("automation-aggregation", "Minutes")
req.Header.Set("x-unknown-future-header", "нечто новое")
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("статус = %d", rec.Code)
}
d, err := st.LastDelivery(t.Context())
if err != nil {
t.Fatalf("LastDelivery: %v", err)
}
var got map[string][]string
if err := json.Unmarshal([]byte(d.Headers), &got); err != nil {
t.Fatalf("разбор headers: %v (%q)", err, d.Headers)
}
for name, want := range map[string]string{
"Automation-Id": "BC99C8A3-8BE7-4519-B545-F3ED6212008E",
"Automation-Aggregation": "Minutes",
"X-Unknown-Future-Header": "нечто новое",
} {
if len(got[name]) != 1 || got[name][0] != want {
t.Errorf("заголовок %s = %v, ожидалось [%q]", name, got[name], want)
}
}
if got["Host"] == nil {
t.Error("Host не сохранён")
}
}
// Секреты в БД не попадают — ни по имени заголовка, ни по совпадению значения
// с настроенным токеном (токен можно положить в заголовок с любым именем).
func TestIngestRedactsSecretHeaders(t *testing.T) {
const token = "очень-секретный-токен"
h, st := newAPI(t, []string{token})
req := httptest.NewRequest(http.MethodPost, "/api/v1/ingest", strings.NewReader(samplePayload))
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("X-Api-Key", "ключ-метабазы")
req.Header.Set("X-Custom-Auth", token)
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("статус = %d, тело %s", rec.Code, rec.Body.String())
}
d, err := st.LastDelivery(t.Context())
if err != nil {
t.Fatalf("LastDelivery: %v", err)
}
if strings.Contains(d.Headers, token) {
t.Errorf("токен утёк в headers: %s", d.Headers)
}
if strings.Contains(d.Headers, "ключ-метабазы") {
t.Errorf("значение X-Api-Key утекло: %s", d.Headers)
}
}
func TestIngestAcceptsGzip(t *testing.T) {
h, st := newAPI(t, nil)
var buf bytes.Buffer
gz := gzip.NewWriter(&buf)
if _, err := gz.Write([]byte(samplePayload)); err != nil {
t.Fatalf("подготовка gzip: %v", err)
}
if err := gz.Close(); err != nil {
t.Fatalf("подготовка gzip: %v", err)
}
req := httptest.NewRequest(http.MethodPost, "/api/v1/ingest", &buf)
req.Header.Set("Content-Encoding", "gzip")
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("статус = %d, тело = %s", rec.Code, rec.Body.String())
}
if n, err := st.CountDeliveries(t.Context()); err != nil || n != 1 {
t.Errorf("доставок = %d (err=%v), ожидалась 1", n, err)
}
}
// Битое тело — отказ доставки, о нём отправителю сообщаем.
func TestIngestRejectsMalformed(t *testing.T) {
h, st := newAPI(t, nil)
req := httptest.NewRequest(http.MethodPost, "/api/v1/ingest", strings.NewReader(`{"data":{"metrics":[`))
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusBadRequest {
t.Fatalf("статус = %d, ожидался 400", rec.Code)
}
if n, _ := st.CountDeliveries(t.Context()); n != 0 { //nolint:errcheck // проверяем только счётчик
t.Errorf("доставок = %d, ожидалось 0", n)
}
}
func TestIngestRejectsOversizedBody(t *testing.T) {
h, _ := newAPI(t, nil)
// Лимит в тесте — 1 МиБ (см. newAPI), шлём заведомо больше.
big := `{"data":{"note":"` + strings.Repeat("x", 2<<20) + `"}}`
req := httptest.NewRequest(http.MethodPost, "/api/v1/ingest", strings.NewReader(big))
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusRequestEntityTooLarge {
t.Fatalf("статус = %d, ожидался 413", rec.Code)
}
}
func TestIngestChecksTokenWhenConfigured(t *testing.T) {
h, _ := newAPI(t, []string{"верный-токен"})
cases := map[string]struct {
header string
want int
}{
"без токена": {"", http.StatusUnauthorized},
"чужой токен": {"Bearer чужой", http.StatusUnauthorized},
"без префикса": {"верный-токен", http.StatusUnauthorized},
"верный токен": {"Bearer верный-токен", http.StatusOK},
}
for name, tc := range cases {
t.Run(name, func(t *testing.T) {
req := httptest.NewRequest(http.MethodPost, "/api/v1/ingest", strings.NewReader(samplePayload))
if tc.header != "" {
req.Header.Set("Authorization", tc.header)
}
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != tc.want {
t.Errorf("статус = %d, ожидался %d", rec.Code, tc.want)
}
})
}
}
func TestHealthz(t *testing.T) {
h, _ := newAPI(t, nil)
rec := httptest.NewRecorder()
h.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/healthz", nil))
if rec.Code != http.StatusOK {
t.Fatalf("статус = %d", rec.Code)
}
}
func newAPI(t *testing.T, writeTokens []string) (http.Handler, *store.Store) {
t.Helper()
dir := t.TempDir()
st, err := store.Open(filepath.Join(dir, "healthlog.db"))
if err != nil {
t.Fatalf("store.Open: %v", err)
}
t.Cleanup(func() { _ = st.Close() })
arch, err := archive.New(filepath.Join(dir, "raw"))
if err != nil {
t.Fatalf("archive.New: %v", err)
}
log := slog.New(slog.DiscardHandler)
h := httpapi.New(httpapi.Options{
Ingest: ingest.New(arch, st, log),
Log: log,
WriteTokens: writeTokens,
MaxBodyMB: 1,
})
return h, st
}
+139
View File
@@ -0,0 +1,139 @@
package httpapi
import (
"compress/gzip"
"errors"
"io"
"net/http"
"strings"
"git.vakhrushev.me/av/healthlog/internal/ingest"
)
// ingestResponse — что видит Health Auto Export в ответ на принятый пакет.
type ingestResponse struct {
DeliveryID string `json:"delivery_id"`
Bytes int64 `json:"bytes"`
SHA256 string `json:"sha256"`
}
// handleIngest принимает пакет от Health Auto Export.
//
// Код ответа отражает ДОСТАВКУ, а не разбор: 200 означает «тело сохранено в
// архив», и этого достаточно, потому что разобрать сохранённое можно всегда.
func (a *api) handleIngest(w http.ResponseWriter, r *http.Request) {
body, err := readBody(w, r, a.maxBody)
if err != nil {
writeReadError(w, err)
return
}
res, err := a.ingest.Accept(r.Context(), body, a.metaFromHeaders(r))
if err != nil {
// Ошибку уже залогировал use-case; транспорт только переводит её в ответ.
if errors.Is(err, ingest.ErrMalformed) {
writeError(w, http.StatusBadRequest, "некорректный формат пакета")
return
}
writeError(w, http.StatusInternalServerError, "внутренняя ошибка")
return
}
writeJSON(w, http.StatusOK, ingestResponse{
DeliveryID: res.DeliveryID,
Bytes: res.Bytes,
SHA256: res.SHA256,
})
}
// metaFromHeaders достаёт то, что автоматизация Health Auto Export
// рассказывает о себе своими заголовками. Именованные поля — те, по которым
// ходят запросы; полный набор кладётся рядом, потому что документация HAE
// заведомо неполна.
func (a *api) metaFromHeaders(r *http.Request) ingest.Meta {
return ingest.Meta{
AutomationName: r.Header.Get("automation-name"),
AutomationID: r.Header.Get("automation-id"),
Aggregation: r.Header.Get("automation-aggregation"),
Period: r.Header.Get("automation-period"),
SessionID: r.Header.Get("session-id"),
Headers: a.safeHeaders(r),
}
}
// secretHeaders — заголовки, которые не сохраняем никогда, как бы ни был
// настроен отправитель.
var secretHeaders = map[string]bool{
"authorization": true,
"proxy-authorization": true,
"cookie": true,
"x-api-key": true,
"api-key": true,
"x-auth-token": true,
}
// redacted подменяет значение, которое оказалось секретом.
const redacted = "[redacted]"
// safeHeaders собирает заголовки запроса для хранения, вычищая секреты.
//
// Двойная защита: имя из чёрного списка вырезается всегда, а любое значение,
// совпавшее с настроенным токеном, подменяется — токен можно положить в
// заголовок с произвольным именем, и угадать его мы не можем.
func (a *api) safeHeaders(r *http.Request) map[string][]string {
out := make(map[string][]string, len(r.Header)+1)
for name, values := range r.Header {
if secretHeaders[strings.ToLower(name)] {
out[name] = []string{redacted}
continue
}
safe := make([]string, len(values))
for i, v := range values {
if tokenAllowed(v, a.writeTokens) || tokenAllowed(strings.TrimPrefix(v, "Bearer "), a.writeTokens) {
safe[i] = redacted
continue
}
safe[i] = v
}
out[name] = safe
}
// Host в r.Header не попадает — а знать, на какое имя пришёл запрос,
// полезно, когда перед сервисом появится прокси.
if r.Host != "" {
out["Host"] = []string{r.Host}
}
return out
}
func readBody(w http.ResponseWriter, r *http.Request, maxBody int64) ([]byte, error) {
var src io.Reader = http.MaxBytesReader(w, r.Body, maxBody)
if strings.Contains(strings.ToLower(r.Header.Get("Content-Encoding")), "gzip") {
gz, err := gzip.NewReader(src)
if err != nil {
return nil, errBadGzip
}
defer func() { _ = gz.Close() }()
src = gz
}
body, err := io.ReadAll(src)
if err != nil {
return nil, err //nolint:wrapcheck // разбирается writeReadError по типу
}
return body, nil
}
var errBadGzip = errors.New("тело не распаковывается как gzip")
func writeReadError(w http.ResponseWriter, err error) {
var tooLarge *http.MaxBytesError
switch {
case errors.As(err, &tooLarge):
writeError(w, http.StatusRequestEntityTooLarge, "пакет больше допустимого размера")
case errors.Is(err, errBadGzip):
writeError(w, http.StatusBadRequest, "тело не распаковывается как gzip")
default:
writeError(w, http.StatusBadRequest, "не удалось прочитать тело запроса")
}
}
+35
View File
@@ -0,0 +1,35 @@
// Package ident — единая точка генерации и разбора идентификаторов (ULID).
//
// ULID сортируется по времени создания, компактен в логах и URL (без дефисов)
// и глобально уникален across сущностей: поиск по голому id находит все
// упоминания. Канонический вид — нижний регистр.
package ident
import (
"errors"
"fmt"
"strings"
"github.com/oklog/ulid/v2"
)
// ErrInvalid — синтаксически некорректный идентификатор. Такой id трактуем
// как несуществующую сущность, без похода в БД.
var ErrInvalid = errors.New("некорректный идентификатор")
// NewID возвращает новый ULID в нижнем регистре.
func NewID() string {
return strings.ToLower(ulid.Make().String())
}
// Parse валидирует внешний идентификатор и приводит его к каноническому виду.
// Вызывается на входных границах (HTTP, CLI) до запроса к БД: сравнение строк
// в SQLite побайтовое, а base32 ULID при декодировании нечувствителен к
// регистру.
func Parse(s string) (string, error) {
id, err := ulid.ParseStrict(strings.ToUpper(strings.TrimSpace(s)))
if err != nil {
return "", fmt.Errorf("%w: %q", ErrInvalid, s)
}
return strings.ToLower(id.String()), nil
}
+58
View File
@@ -0,0 +1,58 @@
package ident_test
import (
"errors"
"strings"
"testing"
"git.vakhrushev.me/av/healthlog/internal/ident"
)
func TestNewIDIsLowercaseULID(t *testing.T) {
id := ident.NewID()
if len(id) != 26 {
t.Errorf("длина = %d, ожидалось 26: %q", len(id), id)
}
if id != strings.ToLower(id) {
t.Errorf("id не в нижнем регистре: %q", id)
}
if ident.NewID() == id {
t.Error("два вызова вернули один и тот же id")
}
}
// Идентификаторы сортируемы по времени создания — на этом держится
// хронология выборок и корреляция в логах.
func TestNewIDSortsByCreationTime(t *testing.T) {
prev := ident.NewID()
for range 10 {
next := ident.NewID()
if next <= prev {
t.Fatalf("порядок нарушен: %q после %q", next, prev)
}
prev = next
}
}
// Внешний id приводится к каноническому виду: сравнение строк в SQLite
// побайтовое, а base32 ULID при декодировании нечувствителен к регистру.
func TestParseNormalizesCase(t *testing.T) {
id := ident.NewID()
got, err := ident.Parse(strings.ToUpper(id))
if err != nil {
t.Fatalf("Parse: %v", err)
}
if got != id {
t.Errorf("Parse вернул %q, ожидалось %q", got, id)
}
}
func TestParseRejectsGarbage(t *testing.T) {
for _, s := range []string{"", "не-id", "01k1", strings.Repeat("z", 26)} {
if _, err := ident.Parse(s); !errors.Is(err, ident.ErrInvalid) {
t.Errorf("Parse(%q): ошибка = %v, ожидалась ErrInvalid", s, err)
}
}
}
+162
View File
@@ -0,0 +1,162 @@
// Package ingest — use-case приёма пакета от Health Auto Export, общий для
// HTTP и (позже) CLI-импорта.
package ingest
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"log/slog"
"git.vakhrushev.me/av/healthlog/internal/archive"
"git.vakhrushev.me/av/healthlog/internal/ident"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// ErrMalformed — тело не разбирается как JSON ожидаемой формы. Это отказ
// доставки (обрыв, обрезанное тело), о нём отправителю сообщаем ошибкой.
// Незнакомое СОДЕРЖИМОЕ (новая метрика, новая форма точки) сюда не относится:
// такой пакет принимается, сохраняется и доразбирается позже.
var ErrMalformed = errors.New("некорректный формат пакета")
// Meta — что рассказала о себе автоматизация Health Auto Export (свои
// заголовки запроса).
type Meta struct {
AutomationName string
AutomationID string
Aggregation string
Period string
SessionID string
// Headers — весь набор заголовков запроса, уже без секретов (вырезает
// транспорт: только он знает список токенов). Документация HAE называет
// пять заголовков, но что приложение шлёт на самом деле — неизвестно, и
// именно из незадокументированного вышли самые полезные находки.
Headers map[string][]string
}
// Result — итог приёма.
type Result struct {
DeliveryID string
Bytes int64
SHA256 string
RawPath string
}
// Service принимает пакеты: сохраняет тело в архив и учитывает доставку.
type Service struct {
arch *archive.Archive
store *store.Store
log *slog.Logger
}
// New собирает use-case приёма.
func New(arch *archive.Archive, st *store.Store, log *slog.Logger) *Service {
return &Service{arch: arch, store: st, log: log.With("capability", "ingest")}
}
// Accept принимает тело пакета: проверяет форму, кладёт в сырой архив и
// заводит запись о доставке.
//
// Порядок важен: сначала тело оказывается на диске, и только потом появляется
// учётная запись. Обратный порядок дал бы учтённую доставку без данных.
//
// Это единственный логирующий чекпоинт приёма — транспорт исход не логирует.
func (s *Service) Accept(ctx context.Context, body []byte, meta Meta) (Result, error) {
if err := checkEnvelope(body); err != nil {
// Некорректный ввод от отправителя — норма жизни, команде разбирать
// нечего: DEBUG, а не ERROR.
s.log.DebugContext(ctx, "delivery rejected", "error", err, "bytes", len(body))
return Result{}, err
}
sum := sha256.Sum256(body)
res := Result{
DeliveryID: ident.NewID(),
Bytes: int64(len(body)),
SHA256: hex.EncodeToString(sum[:]),
}
receivedAt := store.Now()
rawPath, err := s.arch.Write(res.DeliveryID, receivedAt, body)
if err != nil {
s.log.ErrorContext(ctx, "delivery failed", "error", err, "delivery_id", res.DeliveryID)
return Result{}, fmt.Errorf("archive body: %w", err)
}
res.RawPath = rawPath
err = s.store.CreateDelivery(ctx, store.Delivery{
ID: res.DeliveryID,
ReceivedAt: receivedAt,
Headers: encodeHeaders(meta.Headers),
AutomationName: meta.AutomationName,
AutomationID: meta.AutomationID,
Aggregation: meta.Aggregation,
Period: meta.Period,
SessionID: meta.SessionID,
Bytes: res.Bytes,
SHA256: res.SHA256,
RawPath: rawPath,
ParseStatus: store.ParsePending,
})
if err != nil {
// Тело уже на диске — данные не потеряны, но учёта нет. Разбор архива
// на следующем шаге проекта такую доставку подберёт.
s.log.ErrorContext(ctx, "delivery failed", "error", err, "delivery_id", res.DeliveryID, "raw_path", rawPath)
return Result{}, fmt.Errorf("record delivery: %w", err)
}
s.log.InfoContext(ctx, "delivery accepted",
"delivery_id", res.DeliveryID,
"bytes", res.Bytes,
"raw_path", rawPath,
"automation_name", meta.AutomationName,
"aggregation", meta.Aggregation,
"period", meta.Period)
return res, nil
}
// encodeHeaders сериализует заголовки для хранения. Ключи json.Marshal
// сортирует сам, поэтому запись стабильна и её удобно сравнивать между
// доставками. Сбой сериализации не должен ронять приём: заголовки —
// вспомогательные сведения, а не данные, ради которых всё затевалось.
func encodeHeaders(h map[string][]string) string {
if len(h) == 0 {
return "{}"
}
b, err := json.Marshal(h)
if err != nil {
return "{}"
}
return string(b)
}
// envelope — минимальная форма пакета Health Auto Export.
type envelope struct {
Data json.RawMessage `json:"data"`
}
// checkEnvelope проверяет ровно то, что делает пакет доставленным: это JSON,
// в нём есть объект data. Содержимое data не трогаем — его форму определяет
// Apple, и незнакомая форма не повод отвергать доставку.
func checkEnvelope(body []byte) error {
if len(bytes.TrimSpace(body)) == 0 {
return fmt.Errorf("%w: пустое тело", ErrMalformed)
}
var env envelope
if err := json.Unmarshal(body, &env); err != nil {
return fmt.Errorf("%w: %v", ErrMalformed, err) //nolint:errorlint // причину наружу не раскрываем, она уходит в лог
}
if len(env.Data) == 0 {
return fmt.Errorf("%w: нет объекта data", ErrMalformed)
}
if trimmed := bytes.TrimSpace(env.Data); len(trimmed) == 0 || trimmed[0] != '{' {
return fmt.Errorf("%w: data не объект", ErrMalformed)
}
return nil
}
+141
View File
@@ -0,0 +1,141 @@
package ingest_test
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"io"
"log/slog"
"path/filepath"
"testing"
"git.vakhrushev.me/av/healthlog/internal/archive"
"git.vakhrushev.me/av/healthlog/internal/ingest"
"git.vakhrushev.me/av/healthlog/internal/store"
)
func TestAcceptStoresBodyVerbatim(t *testing.T) {
svc, arch, st := newService(t)
// Незнакомая метрика и неизвестная форма точки — не повод отвергать пакет.
body := []byte(`{"data":{"metrics":[{"name":"новая_метрика","units":"?","data":[{"чего_то":1,"date":"2026-07-31 12:00:00 +0300"}]}]}}`)
res, err := svc.Accept(context.Background(), body, ingest.Meta{
AutomationName: "всё подряд",
Aggregation: "None",
Period: "Since Last Sync",
})
if err != nil {
t.Fatalf("Accept: %v", err)
}
if res.Bytes != int64(len(body)) {
t.Errorf("bytes = %d, ожидалось %d", res.Bytes, len(body))
}
sum := sha256.Sum256(body)
if res.SHA256 != hex.EncodeToString(sum[:]) {
t.Errorf("sha256 = %q", res.SHA256)
}
// Тело в архиве — байт в байт то, что прислали.
r, err := arch.Open(res.RawPath)
if err != nil {
t.Fatalf("Open: %v", err)
}
defer func() { _ = r.Close() }()
got, err := io.ReadAll(r)
if err != nil {
t.Fatalf("чтение архива: %v", err)
}
if string(got) != string(body) {
t.Errorf("в архиве %q, прислано %q", got, body)
}
n, err := st.CountDeliveries(context.Background())
if err != nil {
t.Fatalf("CountDeliveries: %v", err)
}
if n != 1 {
t.Errorf("доставок в БД: %d, ожидалась 1", n)
}
}
// Повтор того же тела пока принимается — отсев идентичных доставок отложен
// (docs/plan.md). Проверяем, что повтор не ломается и не затирает первую.
func TestAcceptAllowsRepeatedBody(t *testing.T) {
svc, _, st := newService(t)
body := []byte(`{"data":{"metrics":[]}}`)
first, err := svc.Accept(context.Background(), body, ingest.Meta{})
if err != nil {
t.Fatalf("первый Accept: %v", err)
}
second, err := svc.Accept(context.Background(), body, ingest.Meta{})
if err != nil {
t.Fatalf("второй Accept: %v", err)
}
if first.DeliveryID == second.DeliveryID {
t.Error("у двух доставок совпали идентификаторы")
}
if first.SHA256 != second.SHA256 {
t.Error("у одинаковых тел разошёлся sha256")
}
n, err := st.CountDeliveries(context.Background())
if err != nil {
t.Fatalf("CountDeliveries: %v", err)
}
if n != 2 {
t.Errorf("доставок в БД: %d, ожидалось 2", n)
}
}
func TestAcceptRejectsMalformed(t *testing.T) {
cases := map[string]string{
"пустое тело": ``,
"не JSON": `не json`,
"обрезанный JSON": `{"data":{"metrics":[`,
"нет data": `{"metrics":[]}`,
"data не объект": `{"data":[]}`,
}
for name, body := range cases {
t.Run(name, func(t *testing.T) {
svc, _, st := newService(t)
_, err := svc.Accept(context.Background(), []byte(body), ingest.Meta{})
if !errors.Is(err, ingest.ErrMalformed) {
t.Fatalf("ошибка = %v, ожидалась ErrMalformed", err)
}
// Отвергнутая доставка не оставляет следа в учёте.
n, err := st.CountDeliveries(context.Background())
if err != nil {
t.Fatalf("CountDeliveries: %v", err)
}
if n != 0 {
t.Errorf("доставок в БД: %d, ожидалось 0", n)
}
})
}
}
func newService(t *testing.T) (*ingest.Service, *archive.Archive, *store.Store) {
t.Helper()
dir := t.TempDir()
st, err := store.Open(filepath.Join(dir, "healthlog.db"))
if err != nil {
t.Fatalf("store.Open: %v", err)
}
t.Cleanup(func() { _ = st.Close() })
arch, err := archive.New(filepath.Join(dir, "raw"))
if err != nil {
t.Fatalf("archive.New: %v", err)
}
log := slog.New(slog.DiscardHandler)
return ingest.New(arch, st, log), arch, st
}
+52
View File
@@ -0,0 +1,52 @@
// Package logging собирает slog-логгер по настройкам из конфига.
package logging
import (
"log/slog"
"os"
"strings"
"time"
)
// New возвращает slog-логгер с указанным уровнем и форматом ("json"|"text").
func New(level, format string) *slog.Logger {
opts := &slog.HandlerOptions{Level: parseLevel(level), ReplaceAttr: utcTime}
var handler slog.Handler
if strings.EqualFold(format, "text") {
handler = slog.NewTextHandler(os.Stdout, opts)
} else {
handler = slog.NewJSONHandler(os.Stdout, opts)
}
return slog.New(handler)
}
// NewStderr — JSON-логгер в stderr (UTC) для фатальных ошибок старта, когда
// основной логгер ещё не собран (конфиг не прочитан).
func NewStderr() *slog.Logger {
return slog.New(slog.NewJSONHandler(os.Stderr, &slog.HandlerOptions{ReplaceAttr: utcTime}))
}
// utcTime приводит метку времени записи к UTC: однозначный порядок событий и
// лексикографическая сортировка.
func utcTime(groups []string, a slog.Attr) slog.Attr {
if len(groups) == 0 && a.Key == slog.TimeKey {
if t, ok := a.Value.Any().(time.Time); ok {
a.Value = slog.TimeValue(t.UTC())
}
}
return a
}
func parseLevel(level string) slog.Level {
switch strings.ToLower(level) {
case "debug":
return slog.LevelDebug
case "warn", "warning":
return slog.LevelWarn
case "error":
return slog.LevelError
default:
return slog.LevelInfo
}
}
+97
View File
@@ -0,0 +1,97 @@
package store
import (
"context"
"database/sql"
"errors"
"fmt"
"time"
)
// Статусы разбора доставки.
const (
// ParsePending — тело сохранено, разбора ещё не было.
ParsePending = "pending"
)
// Delivery — учётная запись одного принятого пакета.
type Delivery struct {
ID string
ReceivedAt time.Time
AutomationName string
AutomationID string
Aggregation string
Period string
SessionID string
Bytes int64
SHA256 string
RawPath string
ParseStatus string
Points int64
// Headers — весь набор заголовков запроса как JSON-объект (имя → массив
// значений), уже без секретов. Именованные поля выше дублируют часть из
// них: по ним ходят запросы, а Headers хранит всё остальное на будущее.
Headers string
}
// CreateDelivery записывает факт приёма пакета.
func (s *Store) CreateDelivery(ctx context.Context, d Delivery) error {
const q = `
INSERT INTO delivery (id, received_at, automation_name, automation_id,
aggregation, period, session_id, bytes, sha256,
raw_path, parse_status, points, headers)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`
headers := d.Headers
if headers == "" {
headers = "{}"
}
_, err := s.db.ExecContext(ctx, q,
d.ID, FormatTime(d.ReceivedAt), d.AutomationName, d.AutomationID,
d.Aggregation, d.Period, d.SessionID, d.Bytes, d.SHA256,
d.RawPath, d.ParseStatus, d.Points, headers)
if err != nil {
return fmt.Errorf("insert delivery: %w", err)
}
return nil
}
// LastDelivery возвращает последнюю по времени приёма доставку.
// Возвращает ErrNotFound, если доставок ещё не было.
func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) {
const q = `
SELECT id, received_at, automation_name, automation_id, aggregation,
period, session_id, bytes, sha256, raw_path, parse_status,
points, headers
FROM delivery ORDER BY received_at DESC, id DESC LIMIT 1`
var d Delivery
var receivedAt string
err := s.db.QueryRowxContext(ctx, q).Scan(
&d.ID, &receivedAt, &d.AutomationName, &d.AutomationID, &d.Aggregation,
&d.Period, &d.SessionID, &d.Bytes, &d.SHA256, &d.RawPath,
&d.ParseStatus, &d.Points, &d.Headers)
if errors.Is(err, sql.ErrNoRows) {
return Delivery{}, ErrNotFound
}
if err != nil {
return Delivery{}, fmt.Errorf("select last delivery: %w", err)
}
d.ReceivedAt, err = ParseTime(receivedAt)
if err != nil {
return Delivery{}, err
}
return d, nil
}
// CountDeliveries возвращает число принятых пакетов. Нужно для healthcheck и
// быстрой проверки «данные вообще идут».
func (s *Store) CountDeliveries(ctx context.Context) (int64, error) {
var n int64
if err := s.db.GetContext(ctx, &n, `SELECT count(*) FROM delivery`); err != nil {
return 0, fmt.Errorf("count deliveries: %w", err)
}
return n, nil
}
+8
View File
@@ -0,0 +1,8 @@
package store
import "errors"
// ErrNotFound — записи нет. Граничную ошибку драйвера (sql.ErrNoRows)
// транслируем в доменную здесь же, у источника, чтобы выше по коду не торчал
// database/sql.
var ErrNotFound = errors.New("запись не найдена")
@@ -0,0 +1,26 @@
-- +goose Up
-- Доставка — один принятый пакет от Health Auto Export. Тело лежит в сыром
-- архиве по raw_path; здесь только учёт и метаданные автоматизации.
CREATE TABLE delivery (
id TEXT PRIMARY KEY,
received_at TEXT NOT NULL,
automation_name TEXT NOT NULL DEFAULT '',
automation_id TEXT NOT NULL DEFAULT '',
aggregation TEXT NOT NULL DEFAULT '',
period TEXT NOT NULL DEFAULT '',
session_id TEXT NOT NULL DEFAULT '',
bytes INTEGER NOT NULL,
sha256 TEXT NOT NULL,
raw_path TEXT NOT NULL,
parse_status TEXT NOT NULL,
points INTEGER NOT NULL DEFAULT 0
);
CREATE INDEX delivery_received_at ON delivery (received_at);
-- Хеш тела: пока только учёт. Отсев идентичных повторов включим, когда
-- увидим на живых данных, как часто HAE присылает байт-в-байт одно и то же.
CREATE INDEX delivery_sha256 ON delivery (sha256);
-- +goose Down
DROP TABLE delivery;
@@ -0,0 +1,9 @@
-- +goose Up
-- Полный набор заголовков запроса (JSON: имя → массив значений), кроме
-- несущих секреты. Документация HAE называет пять заголовков, но что
-- приложение шлёт на самом деле — неизвестно, а именно из незадокументированных
-- полей уже вышли самые полезные находки (см. docs/local-research.md).
ALTER TABLE delivery ADD COLUMN headers TEXT NOT NULL DEFAULT '{}';
-- +goose Down
ALTER TABLE delivery DROP COLUMN headers;
+84
View File
@@ -0,0 +1,84 @@
// Package store — SQLite-витрина healthlog: доставки и (позже) разобранные
// точки. Витрина производна от сырого архива и пересобирается из него.
package store
import (
"embed"
"fmt"
"net/url"
"time"
"github.com/jmoiron/sqlx"
"github.com/pressly/goose/v3"
_ "modernc.org/sqlite" // чистый Go-драйвер SQLite, без cgo
)
//go:embed migrations/*.sql
var migrationsFS embed.FS
// Store — доступ к витрине.
type Store struct {
db *sqlx.DB
}
// Open открывает БД по пути и накатывает миграции.
func Open(dbPath string) (*Store, error) {
db, err := sqlx.Connect("sqlite", dsn(dbPath))
if err != nil {
return nil, fmt.Errorf("open sqlite %q: %w", dbPath, err)
}
if err := migrate(db); err != nil {
_ = db.Close()
return nil, err
}
return &Store{db: db}, nil
}
// Close закрывает соединение с БД.
func (s *Store) Close() error {
if err := s.db.Close(); err != nil {
return fmt.Errorf("close sqlite: %w", err)
}
return nil
}
// dsn собирает строку подключения: WAL для параллельного чтения во время
// записи, busy_timeout — чтобы конкурентная запись ждала, а не падала.
func dsn(path string) string {
q := url.Values{}
q.Add("_pragma", "journal_mode(WAL)")
q.Add("_pragma", "busy_timeout(5000)")
q.Add("_pragma", "foreign_keys(on)")
return "file:" + path + "?" + q.Encode()
}
func migrate(db *sqlx.DB) error {
goose.SetBaseFS(migrationsFS)
goose.SetLogger(goose.NopLogger())
if err := goose.SetDialect("sqlite3"); err != nil {
return fmt.Errorf("goose dialect: %w", err)
}
if err := goose.Up(db.DB, "migrations"); err != nil {
return fmt.Errorf("goose up: %w", err)
}
return nil
}
// Now — единая точка генерации времени: UTC, секундная точность.
// Секунды дают фиксированную ширину RFC 3339, а значит лексикографическая
// сортировка TEXT совпадает с хронологией.
func Now() time.Time { return time.Now().UTC().Truncate(time.Second) }
// FormatTime приводит время к каноническому виду хранения: RFC 3339, UTC.
func FormatTime(t time.Time) string { return t.UTC().Format(time.RFC3339) }
// ParseTime разбирает метку времени из БД.
func ParseTime(s string) (time.Time, error) {
t, err := time.Parse(time.RFC3339, s)
if err != nil {
return time.Time{}, fmt.Errorf("parse time %q: %w", s, err)
}
return t.UTC(), nil
}