From 41f2f05f2fe9079c04bd0e4ad9acd287987919de Mon Sep 17 00:00:00 2001 From: Anton Vakhrushev Date: Sat, 1 Aug 2026 12:37:17 +0300 Subject: [PATCH] =?UTF-8?q?=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=20=D0=BF=D1=80=D0=B8=D1=91=D0=BC=20=D0=BF=D0=B0=D0=BA?= =?UTF-8?q?=D0=B5=D1=82=D0=BE=D0=B2=20Health=20Auto=20Export?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - POST /api/v1/ingest: токен, лимит тела, gzip, durable-запись тела в сырой архив, учёт доставки в SQLite - код ответа отражает доставку, а не разбор: битый JSON — 400, непонятое содержимое — 200, данные уже сохранены и доразберутся позже - сохраняется полный набор заголовков запроса с вычисткой секретов по имени и по совпадению значения с токеном --- cmd/healthlog/healthcheck.go | 55 ++++ cmd/healthlog/main.go | 48 ++++ cmd/healthlog/serve.go | 99 +++++++ internal/archive/archive.go | 129 +++++++++ internal/archive/archive_test.go | 92 +++++++ internal/config/config.go | 176 +++++++++++++ internal/config/config_test.go | 111 ++++++++ internal/httpapi/httpapi.go | 138 ++++++++++ internal/httpapi/httpapi_test.go | 244 ++++++++++++++++++ internal/httpapi/ingest.go | 139 ++++++++++ internal/ident/ident.go | 35 +++ internal/ident/ident_test.go | 58 +++++ internal/ingest/ingest.go | 162 ++++++++++++ internal/ingest/ingest_test.go | 141 ++++++++++ internal/logging/logging.go | 52 ++++ internal/store/delivery.go | 97 +++++++ internal/store/errors.go | 8 + internal/store/migrations/00001_delivery.sql | 26 ++ .../migrations/00002_delivery_headers.sql | 9 + internal/store/store.go | 84 ++++++ 20 files changed, 1903 insertions(+) create mode 100644 cmd/healthlog/healthcheck.go create mode 100644 cmd/healthlog/main.go create mode 100644 cmd/healthlog/serve.go create mode 100644 internal/archive/archive.go create mode 100644 internal/archive/archive_test.go create mode 100644 internal/config/config.go create mode 100644 internal/config/config_test.go create mode 100644 internal/httpapi/httpapi.go create mode 100644 internal/httpapi/httpapi_test.go create mode 100644 internal/httpapi/ingest.go create mode 100644 internal/ident/ident.go create mode 100644 internal/ident/ident_test.go create mode 100644 internal/ingest/ingest.go create mode 100644 internal/ingest/ingest_test.go create mode 100644 internal/logging/logging.go create mode 100644 internal/store/delivery.go create mode 100644 internal/store/errors.go create mode 100644 internal/store/migrations/00001_delivery.sql create mode 100644 internal/store/migrations/00002_delivery_headers.sql create mode 100644 internal/store/store.go diff --git a/cmd/healthlog/healthcheck.go b/cmd/healthlog/healthcheck.go new file mode 100644 index 0000000..482b7a6 --- /dev/null +++ b/cmd/healthlog/healthcheck.go @@ -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 +} diff --git a/cmd/healthlog/main.go b/cmd/healthlog/main.go new file mode 100644 index 0000000..c1942c6 --- /dev/null +++ b/cmd/healthlog/main.go @@ -0,0 +1,48 @@ +// Команда healthlog — коллектор данных Apple Health из Health Auto Export. +// +// Подкоманды: +// +// healthlog [serve] --config принимать пакеты (по умолчанию) +// healthlog healthcheck --config

проверить /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) + } +} diff --git a/cmd/healthlog/serve.go b/cmd/healthlog/serve.go new file mode 100644 index 0000000..cb63e0e --- /dev/null +++ b/cmd/healthlog/serve.go @@ -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 +} diff --git a/internal/archive/archive.go b/internal/archive/archive.go new file mode 100644 index 0000000..c9a805c --- /dev/null +++ b/internal/archive/archive.go @@ -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 +} diff --git a/internal/archive/archive_test.go b/internal/archive/archive_test.go new file mode 100644 index 0000000..bb4d6c1 --- /dev/null +++ b/internal/archive/archive_test.go @@ -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 +} diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..3a3d0b3 --- /dev/null +++ b/internal/config/config.go @@ -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) } diff --git a/internal/config/config_test.go b/internal/config/config_test.go new file mode 100644 index 0000000..1a58be4 --- /dev/null +++ b/internal/config/config_test.go @@ -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 +} diff --git a/internal/httpapi/httpapi.go b/internal/httpapi/httpapi.go new file mode 100644 index 0000000..c91848a --- /dev/null +++ b/internal/httpapi/httpapi.go @@ -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}) +} diff --git a/internal/httpapi/httpapi_test.go b/internal/httpapi/httpapi_test.go new file mode 100644 index 0000000..7856bef --- /dev/null +++ b/internal/httpapi/httpapi_test.go @@ -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 +} diff --git a/internal/httpapi/ingest.go b/internal/httpapi/ingest.go new file mode 100644 index 0000000..1d0dcbd --- /dev/null +++ b/internal/httpapi/ingest.go @@ -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, "не удалось прочитать тело запроса") + } +} diff --git a/internal/ident/ident.go b/internal/ident/ident.go new file mode 100644 index 0000000..a0db154 --- /dev/null +++ b/internal/ident/ident.go @@ -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 +} diff --git a/internal/ident/ident_test.go b/internal/ident/ident_test.go new file mode 100644 index 0000000..1eb5eac --- /dev/null +++ b/internal/ident/ident_test.go @@ -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) + } + } +} diff --git a/internal/ingest/ingest.go b/internal/ingest/ingest.go new file mode 100644 index 0000000..b9730ae --- /dev/null +++ b/internal/ingest/ingest.go @@ -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 +} diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go new file mode 100644 index 0000000..2177190 --- /dev/null +++ b/internal/ingest/ingest_test.go @@ -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 +} diff --git a/internal/logging/logging.go b/internal/logging/logging.go new file mode 100644 index 0000000..49fd093 --- /dev/null +++ b/internal/logging/logging.go @@ -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 + } +} diff --git a/internal/store/delivery.go b/internal/store/delivery.go new file mode 100644 index 0000000..a23f17f --- /dev/null +++ b/internal/store/delivery.go @@ -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 +} diff --git a/internal/store/errors.go b/internal/store/errors.go new file mode 100644 index 0000000..b81a88f --- /dev/null +++ b/internal/store/errors.go @@ -0,0 +1,8 @@ +package store + +import "errors" + +// ErrNotFound — записи нет. Граничную ошибку драйвера (sql.ErrNoRows) +// транслируем в доменную здесь же, у источника, чтобы выше по коду не торчал +// database/sql. +var ErrNotFound = errors.New("запись не найдена") diff --git a/internal/store/migrations/00001_delivery.sql b/internal/store/migrations/00001_delivery.sql new file mode 100644 index 0000000..d36c638 --- /dev/null +++ b/internal/store/migrations/00001_delivery.sql @@ -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; diff --git a/internal/store/migrations/00002_delivery_headers.sql b/internal/store/migrations/00002_delivery_headers.sql new file mode 100644 index 0000000..1f38885 --- /dev/null +++ b/internal/store/migrations/00002_delivery_headers.sql @@ -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; diff --git a/internal/store/store.go b/internal/store/store.go new file mode 100644 index 0000000..09f34a9 --- /dev/null +++ b/internal/store/store.go @@ -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 +}