- POST /api/v1/ingest: токен, лимит тела, gzip, durable-запись тела в сырой архив, учёт доставки в SQLite - код ответа отражает доставку, а не разбор: битый JSON — 400, непонятое содержимое — 200, данные уже сохранены и доразберутся позже - сохраняется полный набор заголовков запроса с вычисткой секретов по имени и по совпадению значения с токеном
98 lines
3.2 KiB
Go
98 lines
3.2 KiB
Go
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
|
|
}
|