Add Emby-compatible request logging and ordered playback progress

This commit is contained in:
truewhile
2026-09-12 23:06:44 +08:00
parent 5ddca9a1f6
commit 2989649445
22 changed files with 1105 additions and 77 deletions
+65 -16
View File
@@ -2,6 +2,7 @@ package repository
import (
"context"
"strings"
"gorm.io/gorm"
"gorm.io/gorm/clause"
@@ -13,24 +14,72 @@ import (
// upserts on (UserID, MediaID) so resume always reads the latest position.
type HistoryRepository struct{ db *gorm.DB }
// Upsert atomically inserts/updates the resume position in a single statement,
// relying on the uniq_user_history composite unique index. Concurrent progress
// reports for the same (user, media) can no longer double-insert.
// Upsert atomically inserts/updates the resume position in a single statement.
// Callers without playback-session metadata keep the legacy last-write-wins
// semantics (for example an explicit "mark played" action).
func (r *HistoryRepository) Upsert(ctx context.Context, h *model.PlaybackHistory) error {
return r.db.WithContext(ctx).Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "user_id"}, {Name: "media_id"}},
DoUpdates: clause.Assignments(map[string]any{
"position_ms": h.PositionMs,
// 沿用旧语义:未知时长(0)不覆盖已记录的时长。
"duration_ms": gorm.Expr(
"CASE WHEN ? > 0 THEN ? ELSE playback_histories.duration_ms END",
h.DurationMs, h.DurationMs,
return r.upsert(ctx, h, false)
}
// UpsertProgress writes a progress report only when it is newer than the row
// currently stored. Reports from an older playback session, or an older
// sequence in the same session, are ignored. Reports without session metadata
// fall back to Upsert for compatibility with older clients.
func (r *HistoryRepository) UpsertProgress(ctx context.Context, h *model.PlaybackHistory) error {
h.SessionID = strings.TrimSpace(h.SessionID)
if h.SessionID == "" || h.SessionStartedAtMs <= 0 {
return r.Upsert(ctx, h)
}
return r.upsert(ctx, h, true)
}
func (r *HistoryRepository) upsert(ctx context.Context, h *model.PlaybackHistory, versioned bool) error {
updates := playbackHistoryAssignments(h)
if versioned {
updates["session_id"] = h.SessionID
updates["session_started_at_ms"] = h.SessionStartedAtMs
updates["sequence"] = h.Sequence
}
onConflict := clause.OnConflict{
Columns: []clause.Column{{Name: "user_id"}, {Name: "media_id"}},
DoUpdates: clause.Assignments(updates),
}
if versioned {
onConflict.Where = clause.Where{Exprs: []clause.Expression{
clause.Or(
clause.Lt{
Column: clause.Column{Name: "session_started_at_ms"},
Value: h.SessionStartedAtMs,
},
clause.And(
clause.Eq{
Column: clause.Column{Name: "session_started_at_ms"},
Value: h.SessionStartedAtMs,
},
clause.Lte{
Column: clause.Column{Name: "sequence"},
Value: h.Sequence,
},
),
),
"watched_at": h.WatchedAt,
"completed": h.Completed,
"deleted_at": nil,
}),
}).Create(h).Error
}}
}
return r.db.WithContext(ctx).Clauses(onConflict).Create(h).Error
}
func playbackHistoryAssignments(h *model.PlaybackHistory) map[string]any {
return map[string]any{
"position_ms": h.PositionMs,
// 沿用旧语义:未知时长(0)不覆盖已记录的时长。
"duration_ms": gorm.Expr(
"CASE WHEN ? > 0 THEN ? ELSE playback_histories.duration_ms END",
h.DurationMs, h.DurationMs,
),
"watched_at": h.WatchedAt,
"completed": h.Completed,
"deleted_at": nil,
}
}
// ListByUser returns the most recent history rows for the user.
@@ -77,3 +77,105 @@ func TestHistoryUpsertKeepsDurationWhenUnknown(t *testing.T) {
t.Fatalf("duration_ms=0 upsert must not clear stored duration, got %d", got.DurationMs)
}
}
func TestHistoryUpsertProgressRejectsStaleReports(t *testing.T) {
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err := database.AutoMigrate(db); err != nil {
t.Fatalf("migrate: %v", err)
}
repos := New(db)
ctx := t.Context()
watched := time.Now()
first := &model.PlaybackHistory{
UserID: "u-1",
MediaID: "m-1",
PositionMs: 30_000,
DurationMs: 120_000,
WatchedAt: watched,
SessionID: "session-1",
SessionStartedAtMs: 1_000,
Sequence: 1,
}
if err := repos.History.UpsertProgress(ctx, first); err != nil {
t.Fatalf("first progress: %v", err)
}
newer := &model.PlaybackHistory{
UserID: "u-1",
MediaID: "m-1",
PositionMs: 90_000,
DurationMs: 120_000,
WatchedAt: watched.Add(time.Minute),
Completed: true,
SessionID: "session-1",
SessionStartedAtMs: 1_000,
Sequence: 2,
}
if err := repos.History.UpsertProgress(ctx, newer); err != nil {
t.Fatalf("newer progress: %v", err)
}
stale := &model.PlaybackHistory{
UserID: "u-1",
MediaID: "m-1",
PositionMs: 10_000,
DurationMs: 120_000,
WatchedAt: watched.Add(2 * time.Minute),
SessionID: "session-1",
SessionStartedAtMs: 1_000,
Sequence: 1,
}
if err := repos.History.UpsertProgress(ctx, stale); err != nil {
t.Fatalf("stale progress: %v", err)
}
var got model.PlaybackHistory
if err := db.Where("user_id = ? AND media_id = ?", "u-1", "m-1").First(&got).Error; err != nil {
t.Fatal(err)
}
if got.PositionMs != 90_000 || !got.Completed || got.Sequence != 2 {
t.Fatalf("stale report overwrote newer state: %#v", got)
}
oldSession := &model.PlaybackHistory{
UserID: "u-1",
MediaID: "m-1",
PositionMs: 5_000,
DurationMs: 120_000,
WatchedAt: watched.Add(3 * time.Minute),
SessionID: "session-0",
SessionStartedAtMs: 500,
Sequence: 99,
}
if err := repos.History.UpsertProgress(ctx, oldSession); err != nil {
t.Fatalf("old session progress: %v", err)
}
if err := db.Where("user_id = ? AND media_id = ?", "u-1", "m-1").First(&got).Error; err != nil {
t.Fatal(err)
}
if got.PositionMs != 90_000 || !got.Completed {
t.Fatalf("old session overwrote newer state: %#v", got)
}
restart := &model.PlaybackHistory{
UserID: "u-1",
MediaID: "m-1",
PositionMs: 1_000,
DurationMs: 120_000,
WatchedAt: watched.Add(4 * time.Minute),
SessionID: "session-2",
SessionStartedAtMs: 2_000,
Sequence: 1,
}
if err := repos.History.UpsertProgress(ctx, restart); err != nil {
t.Fatalf("restart progress: %v", err)
}
if err := db.Where("user_id = ? AND media_id = ?", "u-1", "m-1").First(&got).Error; err != nil {
t.Fatal(err)
}
if got.PositionMs != 1_000 || got.Completed || got.SessionID != "session-2" {
t.Fatalf("new playback session did not reset state: %#v", got)
}
}