feat(stats): #10 第②层 —— 统计页按设备过滤(per-device 本地日曲线)
ci-pangolin / Lint — shellcheck (push) Successful in 8s
ci-pangolin / OpenAPI Sync Check (push) Successful in 19s
ci-pangolin / Redline Scan — 脱敏 (UI 文案) (push) Successful in 6s
ci-pangolin / Flutter — analyze + test (push) Failing after 14s
ci-pangolin / Portable SQL — 可移植性 (mysql/sqlite) (push) Successful in 6s
ci-pangolin / Codegen Drift — token 生成物未漂移 (push) Successful in 4s
ci-pangolin / Go — build + test (push) Successful in 11s
ci-pangolin / E2E Smoke — L4 进程级端到端 (push) Successful in 14s
ci-pangolin / Go — integration (mysql/redis testcontainers) (push) Failing after 4m8s
ci-pangolin / Golden — 视觉回归 (components + auth) (push) Successful in 15s
ci-pangolin / Lint — shellcheck (push) Successful in 8s
ci-pangolin / OpenAPI Sync Check (push) Successful in 19s
ci-pangolin / Redline Scan — 脱敏 (UI 文案) (push) Successful in 6s
ci-pangolin / Flutter — analyze + test (push) Failing after 14s
ci-pangolin / Portable SQL — 可移植性 (mysql/sqlite) (push) Successful in 6s
ci-pangolin / Codegen Drift — token 生成物未漂移 (push) Successful in 4s
ci-pangolin / Go — build + test (push) Successful in 11s
ci-pangolin / E2E Smoke — L4 进程级端到端 (push) Successful in 14s
ci-pangolin / Go — integration (mysql/redis testcontainers) (push) Failing after 4m8s
ci-pangolin / Golden — 视觉回归 (components + auth) (push) Successful in 15s
镜像账户级时区曲线到「按设备」: - migration 000018 `usage_device_hourly`(每设备 UTC 小时桶,稀疏)。 - ReportUsage deviceID>0 时加 AccumulateDeviceHourly;NodeStore 接口+impl+mock。 - usage.Store.DeviceHourlyRange(JOIN devices 校验归属+解析 uuid→id);UsageCurve 增 deviceUUID 形参:空=账户级,非空=该设备本地日曲线(分桶逻辑复用)。 - /v1/usage?device=<uuid>;客户端 account_api.usage(device)、usageProvider key 改 记录 (days,device)、stats_page 接 statsDeviceProvider → 选设备即重取该设备曲线。 测试:per-device 曲线隔离+归属校验(A的设备不进B)、UsageCurve(deviceUUID)、 migration v18、客户端 widget(选设备→/v1/usage 带 device=)。全量 go/flutter 绿。 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -132,6 +132,12 @@ func (m *mockNodeStore) AccumulateHourly(_ context.Context, userID int64, hour t
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockNodeStore) AccumulateDeviceHourly(_ context.Context, _, _ int64, _ time.Time,
|
||||
_, _, _ int64,
|
||||
) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockNodeStore) EnsureDeviceDpUUID(_ context.Context, _ int64, _ string) (string, int64, error) {
|
||||
return "", 0, nil
|
||||
}
|
||||
|
||||
@@ -325,6 +325,13 @@ func (h *Handler) ReportUsage(ctx context.Context, req *agentv1.UsageReport) (*a
|
||||
slog.Warn("nodes.Handler.ReportUsage: device accumulate failed",
|
||||
"user_id", userID, "device_id", deviceID, "err", err)
|
||||
}
|
||||
// Per-device hourly (tz-aware per-device display curve, #10②).
|
||||
if err := h.store.AccumulateDeviceHourly(ctx, userID, deviceID, windowEnd,
|
||||
entry.BytesUp, entry.BytesDown, entry.SessionMinutes,
|
||||
); err != nil {
|
||||
slog.Warn("nodes.Handler.ReportUsage: device hourly accumulate failed",
|
||||
"user_id", userID, "device_id", deviceID, "err", err)
|
||||
}
|
||||
// Heartbeat: keep the device's "online" status fresh while it reports.
|
||||
if err := h.store.TouchDeviceLastSeen(ctx, deviceID); err != nil {
|
||||
slog.Warn("nodes.Handler.ReportUsage: touch last_seen failed",
|
||||
|
||||
@@ -108,6 +108,12 @@ type NodeStore interface {
|
||||
AccumulateDeviceUsage(ctx context.Context, userID, deviceID int64, date time.Time,
|
||||
bytesUp, bytesDown int64, minutes int64) error
|
||||
|
||||
// AccumulateDeviceHourly adds bytes/minutes to usage_device_hourly keyed by
|
||||
// (deviceID, UTC epoch hour) — the per-device counterpart of AccumulateHourly,
|
||||
// for timezone-aware per-device day bucketing at query time.
|
||||
AccumulateDeviceHourly(ctx context.Context, userID, deviceID int64, hour time.Time,
|
||||
bytesUp, bytesDown int64, minutes int64) error
|
||||
|
||||
// TouchDeviceLastSeen bumps devices.last_seen to now for an active device,
|
||||
// driving the "online" status (connect + periodic usage reports refresh it).
|
||||
TouchDeviceLastSeen(ctx context.Context, deviceID int64) error
|
||||
@@ -501,3 +507,23 @@ func (s *SQLNodeStore) AccumulateHourly(
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// AccumulateDeviceHourly adds bytes/minutes to usage_device_hourly for
|
||||
// (deviceID, UTC epoch hour). Per-device counterpart of AccumulateHourly.
|
||||
func (s *SQLNodeStore) AccumulateDeviceHourly(
|
||||
ctx context.Context, userID, deviceID int64, hour time.Time,
|
||||
bytesUp, bytesDown int64, minutes int64,
|
||||
) error {
|
||||
h := hour.UTC().Unix() / 3600
|
||||
q := `
|
||||
INSERT INTO usage_device_hourly (user_id, device_id, hour, bytes_up, bytes_down, minutes_used)
|
||||
VALUES (?, ?, ?, ?, ?, ?)
|
||||
` + s.dialect.Upsert([]string{"device_id", "hour"},
|
||||
"bytes_up = bytes_up + EXCLUDED.bytes_up",
|
||||
"bytes_down = bytes_down + EXCLUDED.bytes_down",
|
||||
"minutes_used = minutes_used + EXCLUDED.minutes_used")
|
||||
if _, err := s.db.ExecContext(ctx, q, userID, deviceID, h, bytesUp, bytesDown, minutes); err != nil {
|
||||
return fmt.Errorf("nodes.SQLNodeStore.AccumulateDeviceHourly: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -29,8 +29,8 @@ func TestSQLiteMigrateUpDown(t *testing.T) {
|
||||
if dirty {
|
||||
t.Fatalf("schema dirty after MigrateUp")
|
||||
}
|
||||
if v != 17 {
|
||||
t.Errorf("version = %d, want 17", v)
|
||||
if v != 18 {
|
||||
t.Errorf("version = %d, want 18", v)
|
||||
}
|
||||
|
||||
// 2. Core tables exist.
|
||||
@@ -38,7 +38,7 @@ func TestSQLiteMigrateUpDown(t *testing.T) {
|
||||
"users", "devices", "plans", "subscriptions", "code_batches", "codes",
|
||||
"usage_daily", "audit_log", "providers", "nodes", "node_events",
|
||||
"directory_version", "provision_idempotency", "replacements", "admins",
|
||||
"connect_credentials", "usage_device_daily", "sessions", "usage_hourly",
|
||||
"connect_credentials", "usage_device_daily", "sessions", "usage_hourly", "usage_device_hourly",
|
||||
} {
|
||||
var name string
|
||||
err := db.QueryRow(
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
package store_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/wangjia/pangolin/server/internal/usage"
|
||||
)
|
||||
|
||||
// Per-device curve (#10②): UsageCurve(deviceUUID) must return only that device's
|
||||
// data, account curve (deviceUUID="") the total, and a user must never read
|
||||
// another user's device (JOIN enforces ownership).
|
||||
func TestSQLite_UsageCurve_PerDevice(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
db := openSQLite(t)
|
||||
mustExec := func(q string, args ...any) {
|
||||
if _, err := db.Exec(q, args...); err != nil {
|
||||
t.Fatalf("seed %q: %v", q, err)
|
||||
}
|
||||
}
|
||||
mustExec(`INSERT INTO users (id, uuid, email, pw_hash, dp_uuid, status) VALUES (1,'u1','u1@e','h','dp1','active'),(2,'u2','u2@e','h','dp2','active')`)
|
||||
mustExec(`INSERT INTO devices (id, uuid, user_id, name, platform) VALUES (10,'devA',1,'A','macos'),(20,'devB',1,'B','ios'),(30,'devC',2,'C','windows')`)
|
||||
|
||||
h := time.Now().UTC().Unix()/3600 - 2 // 2h ago, within window
|
||||
mustExec(`INSERT INTO usage_device_hourly (user_id, device_id, hour, bytes_up, bytes_down, minutes_used) VALUES (1,10,?,1000,2000,5),(1,20,?,300,400,2),(2,30,?,999,0,9)`, h, h, h)
|
||||
mustExec(`INSERT INTO usage_hourly (user_id, hour, bytes_up, bytes_down, minutes_used) VALUES (1,?,1300,2400,7)`, h)
|
||||
|
||||
svc := usage.NewService(usage.NewStore(db), nil, nil, time.Hour)
|
||||
sum := func(uid int64, dev string) (up, down uint64, mins int) {
|
||||
pts, e := svc.UsageCurve(ctx, uid, 3, 0, dev)
|
||||
if e != nil {
|
||||
t.Fatalf("UsageCurve(dev=%q): %v", dev, e)
|
||||
}
|
||||
for _, p := range pts {
|
||||
up += p.BytesUp
|
||||
down += p.BytesDown
|
||||
mins += p.MinutesUsed
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
if up, down, m := sum(1, "devA"); up != 1000 || down != 2000 || m != 5 {
|
||||
t.Errorf("devA curve = %d/%d/%d, want 1000/2000/5", up, down, m)
|
||||
}
|
||||
if up, _, _ := sum(1, "devB"); up != 300 {
|
||||
t.Errorf("devB curve up = %d, want 300", up)
|
||||
}
|
||||
if up, _, _ := sum(1, ""); up != 1300 {
|
||||
t.Errorf("account curve up = %d, want 1300", up)
|
||||
}
|
||||
// user 1 querying user 2's device → empty (ownership via JOIN).
|
||||
if up, _, _ := sum(1, "devC"); up != 0 {
|
||||
t.Errorf("cross-user leak: devC visible to user1, up=%d", up)
|
||||
}
|
||||
}
|
||||
@@ -37,7 +37,7 @@ func TestSQLite_UsageCurve_LocalDayBucketing(t *testing.T) {
|
||||
svc := usage.NewService(usage.NewStore(db), nil, nil, time.Hour)
|
||||
|
||||
find := func(days, offset int, date string) (usage.UsagePoint, bool) {
|
||||
pts, e := svc.UsageCurve(ctx, 1, days, offset)
|
||||
pts, e := svc.UsageCurve(ctx, 1, days, offset, "")
|
||||
if e != nil {
|
||||
t.Fatalf("UsageCurve(off=%d): %v", offset, e)
|
||||
}
|
||||
|
||||
@@ -67,7 +67,10 @@ func (h *UsageHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
}
|
||||
|
||||
points, apiErr := h.svc.UsageCurve(r.Context(), userID, days, offsetMin)
|
||||
// device (UUID, optional): when set, return that device's curve only (#10②).
|
||||
deviceUUID := r.URL.Query().Get("device")
|
||||
|
||||
points, apiErr := h.svc.UsageCurve(r.Context(), userID, days, offsetMin, deviceUUID)
|
||||
if apiErr != nil {
|
||||
apierr.WriteJSON(w, http.StatusInternalServerError, apiErr)
|
||||
return
|
||||
|
||||
@@ -53,7 +53,8 @@ type UsagePoint struct {
|
||||
// It reads UTC hourly buckets (usage_hourly) and re-aggregates by local calendar
|
||||
// day, so "today" matches the user's local date regardless of server timezone
|
||||
// (fixes "29号却显示28号数据"). offsetMin out of ±14h falls back to 0 (UTC).
|
||||
func (svc *Service) UsageCurve(ctx context.Context, userID int64, days, offsetMin int) ([]UsagePoint, *apierr.Error) {
|
||||
// deviceUUID == "" → account-level curve; non-empty → that device's curve only.
|
||||
func (svc *Service) UsageCurve(ctx context.Context, userID int64, days, offsetMin int, deviceUUID string) ([]UsagePoint, *apierr.Error) {
|
||||
if days < 1 {
|
||||
days = 7
|
||||
}
|
||||
@@ -73,7 +74,13 @@ func (svc *Service) UsageCurve(ctx context.Context, userID int64, days, offsetMi
|
||||
fromHour := fromLocal.UTC().Unix() / 3600
|
||||
toHour := now.Unix() / 3600
|
||||
|
||||
rows, err := svc.store.HourlyRange(ctx, userID, fromHour, toHour)
|
||||
var rows []HourBucket
|
||||
var err error
|
||||
if deviceUUID == "" {
|
||||
rows, err = svc.store.HourlyRange(ctx, userID, fromHour, toHour)
|
||||
} else {
|
||||
rows, err = svc.store.DeviceHourlyRange(ctx, userID, deviceUUID, fromHour, toHour)
|
||||
}
|
||||
if err != nil {
|
||||
return nil, apierr.ErrInternal
|
||||
}
|
||||
|
||||
@@ -139,6 +139,33 @@ func (s *Store) HourlyRange(ctx context.Context, userID, fromHour, toHour int64)
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// DeviceHourlyRange returns one device's UTC hourly buckets (usage_device_hourly)
|
||||
// for hour in [fromHour, toHour], oldest first. The JOIN both resolves the
|
||||
// client-facing device UUID to its internal id and enforces ownership (user_id),
|
||||
// so a user can never read another user's device curve.
|
||||
func (s *Store) DeviceHourlyRange(ctx context.Context, userID int64, deviceUUID string, fromHour, toHour int64) ([]HourBucket, error) {
|
||||
rows, err := s.db.QueryContext(ctx,
|
||||
`SELECT uh.hour, uh.bytes_up, uh.bytes_down, uh.minutes_used
|
||||
FROM usage_device_hourly uh
|
||||
JOIN devices d ON d.id = uh.device_id
|
||||
WHERE d.user_id = ? AND d.uuid = ? AND uh.hour BETWEEN ? AND ?
|
||||
ORDER BY uh.hour ASC`,
|
||||
userID, deviceUUID, fromHour, toHour)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("store.DeviceHourlyRange: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []HourBucket
|
||||
for rows.Next() {
|
||||
var b HourBucket
|
||||
if err := rows.Scan(&b.Hour, &b.BytesUp, &b.BytesDown, &b.MinutesUsed); err != nil {
|
||||
return nil, fmt.Errorf("store.DeviceHourlyRange scan: %w", err)
|
||||
}
|
||||
out = append(out, b)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// DeviceUsageRow is one device's aggregated usage over a date window, joined to
|
||||
// the devices table for display metadata. It is the per-device ("下分设备")
|
||||
// counterpart of DailyUsage's account rollup.
|
||||
|
||||
@@ -110,6 +110,14 @@ func applySchema(db *sql.DB) error {
|
||||
ad_unlocked_at DATETIME(6) NULL,
|
||||
PRIMARY KEY (user_id, date)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4`,
|
||||
`CREATE TABLE IF NOT EXISTS usage_hourly (
|
||||
user_id BIGINT UNSIGNED NOT NULL,
|
||||
hour BIGINT NOT NULL,
|
||||
bytes_up BIGINT UNSIGNED NOT NULL DEFAULT 0,
|
||||
bytes_down BIGINT UNSIGNED NOT NULL DEFAULT 0,
|
||||
minutes_used INT NOT NULL DEFAULT 0,
|
||||
PRIMARY KEY (user_id, hour)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4`,
|
||||
`INSERT IGNORE INTO plans (code, max_devices, daily_minutes, ad_gate)
|
||||
VALUES ('free',1,10,TRUE), ('pro',5,NULL,FALSE), ('team',10,NULL,FALSE)`,
|
||||
}
|
||||
@@ -445,7 +453,7 @@ func TestIntUsageCurveZeroFill(t *testing.T) {
|
||||
db.Exec(`INSERT INTO usage_hourly (user_id, hour, bytes_up, bytes_down, minutes_used) VALUES (?,?,?,?,?)`,
|
||||
uid, epochHour(today.AddDate(0, 0, -3)), 50, 60, 2)
|
||||
|
||||
points, apiErr := svc.UsageCurve(context.Background(), uid, 7, 0)
|
||||
points, apiErr := svc.UsageCurve(context.Background(), uid, 7, 0, "")
|
||||
if apiErr != nil {
|
||||
t.Fatalf("curve: %v", apiErr)
|
||||
}
|
||||
@@ -486,7 +494,7 @@ func TestIntUsageCurveOnlyCurrentUser(t *testing.T) {
|
||||
db.Exec(`INSERT INTO usage_hourly (user_id, hour, bytes_up, bytes_down, minutes_used) VALUES (?,?,?,?,?)`, uA, h, 111, 0, 1)
|
||||
db.Exec(`INSERT INTO usage_hourly (user_id, hour, bytes_up, bytes_down, minutes_used) VALUES (?,?,?,?,?)`, uB, h, 999, 0, 9)
|
||||
|
||||
points, _ := svc.UsageCurve(context.Background(), uA, 1, 0)
|
||||
points, _ := svc.UsageCurve(context.Background(), uA, 1, 0, "")
|
||||
if len(points) != 1 || points[0].BytesUp != 111 {
|
||||
t.Errorf("user A curve leaked other user data: %+v", points)
|
||||
}
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
DROP TABLE IF EXISTS usage_device_hourly;
|
||||
@@ -0,0 +1,11 @@
|
||||
-- 见 sqlite/000018 注释。
|
||||
CREATE TABLE usage_device_hourly (
|
||||
user_id BIGINT UNSIGNED NOT NULL,
|
||||
device_id BIGINT UNSIGNED NOT NULL,
|
||||
hour BIGINT NOT NULL,
|
||||
bytes_up BIGINT UNSIGNED NOT NULL DEFAULT 0,
|
||||
bytes_down BIGINT UNSIGNED NOT NULL DEFAULT 0,
|
||||
minutes_used INT NOT NULL DEFAULT 0,
|
||||
PRIMARY KEY (device_id, hour),
|
||||
INDEX idx_usage_device_hourly_user (user_id, hour)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
|
||||
@@ -0,0 +1,2 @@
|
||||
DROP INDEX IF EXISTS idx_usage_device_hourly_user;
|
||||
DROP TABLE IF EXISTS usage_device_hourly;
|
||||
@@ -0,0 +1,13 @@
|
||||
-- 每设备 UTC 小时桶,供统计页「按设备」本地日曲线(#10 第②层)。
|
||||
-- 镜像 usage_hourly + device_id(如同 usage_device_daily 之于 usage_daily)。
|
||||
-- hour = unix 纪元小时(UTC);稀疏(空小时不落行)。
|
||||
CREATE TABLE usage_device_hourly (
|
||||
user_id INTEGER NOT NULL,
|
||||
device_id INTEGER NOT NULL,
|
||||
hour INTEGER NOT NULL,
|
||||
bytes_up INTEGER NOT NULL DEFAULT 0,
|
||||
bytes_down INTEGER NOT NULL DEFAULT 0,
|
||||
minutes_used INTEGER NOT NULL DEFAULT 0,
|
||||
PRIMARY KEY (device_id, hour)
|
||||
);
|
||||
CREATE INDEX idx_usage_device_hourly_user ON usage_device_hourly (user_id, hour);
|
||||
Reference in New Issue
Block a user