Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions graphql/schema/schema.graphql
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,17 @@ type Query {
"""Find a performer by ID"""
findPerformer(id: ID!): Performer @hasRole(role: READ)
queryPerformers(input: PerformerQueryInput!): QueryPerformersResultType! @hasRole(role: READ)
"""Incremental feed of performers changed since a keyset cursor"""
performerChangelog(since: Time!, after_id: ID, limit: Int = 5000): [EntityChange!]! @hasRole(role: READ)

#### Studios ####

# studio names should be unique
"""Find a studio by ID or name"""
findStudio(id: ID, name: String): Studio @hasRole(role: READ)
queryStudios(input: StudioQueryInput!): QueryStudiosResultType! @hasRole(role: READ)
"""Incremental feed of studios changed since a keyset cursor"""
studioChangelog(since: Time!, after_id: ID, limit: Int = 5000): [EntityChange!]! @hasRole(role: READ)

#### Tags ####

Expand All @@ -26,6 +30,8 @@ type Query {
"""Find a tag category by ID"""
findTagCategory(id: ID!): TagCategory @hasRole(role: READ)
queryTagCategories: QueryTagCategoriesResultType! @hasRole(role: READ)
"""Incremental feed of tags changed since a keyset cursor"""
tagChangelog(since: Time!, after_id: ID, limit: Int = 5000): [EntityChange!]! @hasRole(role: READ)

#### Scenes ####

Expand All @@ -37,6 +43,8 @@ type Query {
findScenesBySceneFingerprints(fingerprints: [[FingerprintQueryInput!]!]!): [[Scene]!]! @hasRole(role: READ)

queryScenes(input: SceneQueryInput!): QueryScenesResultType! @hasRole(role: READ)
"""Incremental feed of scenes changed since a keyset cursor"""
sceneChangelog(since: Time!, after_id: ID, limit: Int = 5000): [EntityChange!]! @hasRole(role: READ)

"""Find an external site by ID"""
findSite(id: ID!): Site @hasRole(role: READ)
Expand Down
11 changes: 11 additions & 0 deletions graphql/schema/types/misc.graphql
Original file line number Diff line number Diff line change
Expand Up @@ -29,3 +29,14 @@ input URLInput {
url: String!
site_id: ID!
}

# A single entity change reported by the *Changelog feeds. Compact by design:
# clients intersect the ids against their tracked set and hydrate matches via
# the existing find* queries. redirect_to is set when the entity was merged
# into a surviving entity (always with deleted = true).
type EntityChange {
id: ID!
updated_at: Time!
deleted: Boolean!
redirect_to: ID
}
220 changes: 220 additions & 0 deletions internal/api/changelog_integration_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,220 @@
//go:build integration

package api_test

import (
"testing"

"github.com/gofrs/uuid"
"github.com/stashapp/stash-box/internal/models"
"github.com/stretchr/testify/assert"
)

// changelogEpoch is a since value far enough in the past to return everything.
const changelogEpoch = "1970-01-01T00:00:00Z"

type changelogTestRunner struct {
testRunner
}

func createChangelogTestRunner(t *testing.T) *changelogTestRunner {
return &changelogTestRunner{
testRunner: *asAdmin(t),
}
}

func findChange(changes []entityChange, id string) *entityChange {
for i := range changes {
if changes[i].ID == id {
return &changes[i]
}
}
return nil
}

// sceneCursor creates a sentinel scene and returns a keyset cursor positioned at
// it. Scenes created afterwards get strictly greater updated_at (separate
// transactions -> distinct now()), so the cursor cleanly excludes the sentinel
// and isolates the test from data created by earlier tests sharing the DB.
func (s *changelogTestRunner) sceneCursor() (since string, afterID uuid.UUID) {
sentinelScene, err := s.createTestScene(nil)
assert.NoError(s.t, err)

base, err := s.client.changelog("sceneChangelog", changelogEpoch, nil, nil)
assert.NoError(s.t, err)

sentinel := findChange(base, sentinelScene.ID)
assert.NotNil(s.t, sentinel, "sentinel scene should appear in changelog")
if sentinel == nil {
return changelogEpoch, uuid.Nil
}
return sentinel.UpdatedAt, uuid.FromStringOrNil(sentinel.ID)
}

func (s *changelogTestRunner) testSceneChangelogCursor() {
since, afterID := s.sceneCursor()

s1, err := s.createTestScene(nil)
assert.NoError(s.t, err)
s2, err := s.createTestScene(nil)
assert.NoError(s.t, err)

changes, err := s.client.changelog("sceneChangelog", since, &afterID, nil)
assert.NoError(s.t, err)

assert.Nil(s.t, findChange(changes, afterID.String()), "scene at the cursor should be excluded")

c1 := findChange(changes, s1.ID)
c2 := findChange(changes, s2.ID)
assert.NotNil(s.t, c1, "scene created after the cursor should appear")
assert.NotNil(s.t, c2, "scene created after the cursor should appear")
if c1 != nil {
assert.False(s.t, c1.Deleted)
assert.Nil(s.t, c1.RedirectTo)
}
}

// testSceneChangelogPagination walks the feed one row at a time and asserts that
// every created scene is returned exactly once. Completeness + no duplicates is
// a full end-to-end check of the (updated_at, id) keyset: a skip breaks
// completeness, a repeat breaks uniqueness.
func (s *changelogTestRunner) testSceneChangelogPagination() {
since, afterID := s.sceneCursor()

created := map[string]bool{}
for i := 0; i < 3; i++ {
sc, err := s.createTestScene(nil)
assert.NoError(s.t, err)
created[sc.ID] = true
}

limit := 1
seen := map[string]bool{}
for iter := 0; iter < 100; iter++ {
page, err := s.client.changelog("sceneChangelog", since, &afterID, &limit)
assert.NoError(s.t, err)
if len(page) == 0 {
break
}
assert.Len(s.t, page, 1, "limit should cap the page size")
row := page[0]
assert.False(s.t, seen[row.ID], "row returned twice during pagination")
seen[row.ID] = true
since = row.UpdatedAt
afterID = uuid.FromStringOrNil(row.ID)
}

for id := range created {
assert.True(s.t, seen[id], "created scene missing from paginated changelog")
}
}

func (s *changelogTestRunner) testSceneChangelogDeleted() {
since, afterID := s.sceneCursor()

scene, err := s.createTestScene(nil)
assert.NoError(s.t, err)
sceneID := scene.UUID()

editInput := models.EditInput{Operation: models.OperationEnumDestroy, ID: &sceneID}
edit, err := s.createTestSceneEdit(models.OperationEnumDestroy, &models.SceneEditDetailsInput{}, &editInput)
assert.NoError(s.t, err)
_, err = s.approveEdit(edit.ID)
assert.NoError(s.t, err)

changes, err := s.client.changelog("sceneChangelog", since, &afterID, nil)
assert.NoError(s.t, err)

ch := findChange(changes, scene.ID)
assert.NotNil(s.t, ch, "deleted scene should appear in changelog")
if ch != nil {
assert.True(s.t, ch.Deleted, "scene should be marked deleted")
assert.Nil(s.t, ch.RedirectTo, "a plain delete should have no redirect")
}
}

func (s *changelogTestRunner) testSceneChangelogMerge() {
since, afterID := s.sceneCursor()

primary, err := s.createTestScene(nil)
assert.NoError(s.t, err)
source, err := s.createTestScene(nil)
assert.NoError(s.t, err)

primaryID := primary.UUID()
editInput := models.EditInput{
Operation: models.OperationEnumMerge,
ID: &primaryID,
MergeSourceIds: []uuid.UUID{source.UUID()},
}
edit, err := s.createTestSceneEdit(models.OperationEnumMerge, nil, &editInput)
assert.NoError(s.t, err)
_, err = s.approveEdit(edit.ID)
assert.NoError(s.t, err)

changes, err := s.client.changelog("sceneChangelog", since, &afterID, nil)
assert.NoError(s.t, err)

ch := findChange(changes, source.ID)
assert.NotNil(s.t, ch, "merged-away scene should appear in changelog")
if ch != nil {
assert.True(s.t, ch.Deleted, "merged scene should be marked deleted")
assert.NotNil(s.t, ch.RedirectTo, "merged scene should carry a redirect")
if ch.RedirectTo != nil {
assert.Equal(s.t, primary.ID, *ch.RedirectTo, "redirect should point to the surviving scene")
}
}
}

// testEntityChangelogs exercises the performer/studio/tag feeds, which share the
// resolver/query/index pattern with scenes.
func (s *changelogTestRunner) testEntityChangelogs() {
performer, err := s.createTestPerformer(nil)
assert.NoError(s.t, err)
studio, err := s.createTestStudio(nil)
assert.NoError(s.t, err)
tag, err := s.createTestTag(nil)
assert.NoError(s.t, err)

performers, err := s.client.changelog("performerChangelog", changelogEpoch, nil, nil)
assert.NoError(s.t, err)
assert.NotNil(s.t, findChange(performers, performer.ID), "performer should appear in performerChangelog")

studios, err := s.client.changelog("studioChangelog", changelogEpoch, nil, nil)
assert.NoError(s.t, err)
assert.NotNil(s.t, findChange(studios, studio.ID), "studio should appear in studioChangelog")

tags, err := s.client.changelog("tagChangelog", changelogEpoch, nil, nil)
assert.NoError(s.t, err)
assert.NotNil(s.t, findChange(tags, tag.ID), "tag should appear in tagChangelog")
}

func (s *changelogTestRunner) testChangelogRequiresRead() {
none := asNone(s.t)
_, err := none.client.changelog("sceneChangelog", changelogEpoch, nil, nil)
assert.Error(s.t, err, "changelog should require the READ role")
}

func TestSceneChangelogCursor(t *testing.T) {
createChangelogTestRunner(t).testSceneChangelogCursor()
}

func TestSceneChangelogPagination(t *testing.T) {
createChangelogTestRunner(t).testSceneChangelogPagination()
}

func TestSceneChangelogDeleted(t *testing.T) {
createChangelogTestRunner(t).testSceneChangelogDeleted()
}

func TestSceneChangelogMerge(t *testing.T) {
createChangelogTestRunner(t).testSceneChangelogMerge()
}

func TestEntityChangelogs(t *testing.T) {
createChangelogTestRunner(t).testEntityChangelogs()
}

func TestChangelogRequiresRead(t *testing.T) {
createChangelogTestRunner(t).testChangelogRequiresRead()
}
32 changes: 32 additions & 0 deletions internal/api/graphql_client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -307,6 +307,38 @@ type graphqlClient struct {
*client.Client
}

type entityChange struct {
ID string `json:"id"`
UpdatedAt string `json:"updated_at"`
Deleted bool `json:"deleted"`
RedirectTo *string `json:"redirect_to"`
}

// changelog queries one of the *Changelog feeds. field is the query name, e.g.
// "sceneChangelog". afterID/limit may be nil to send null.
func (c *graphqlClient) changelog(field, since string, afterID *uuid.UUID, limit *int) ([]entityChange, error) {
q := `
query Changelog($since: Time!, $after_id: ID, $limit: Int) {
` + field + `(since: $since, after_id: $after_id, limit: $limit) {
id
updated_at
deleted
redirect_to
}
}`

var resp map[string][]entityChange
if err := c.Post(q, &resp,
client.Var("since", since),
client.Var("after_id", afterID),
client.Var("limit", limit),
); err != nil {
return nil, err
}

return resp[field], nil
}

func (c *graphqlClient) createScene(input models.SceneCreateInput) (*sceneOutput, error) {
q := `
mutation SceneCreate($input: SceneCreateInput!) {
Expand Down
47 changes: 47 additions & 0 deletions internal/api/resolver_query_changelog.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
package api

import (
"context"
"time"

"github.com/gofrs/uuid"

"github.com/stashapp/stash-box/internal/models"
)

const defaultChangelogLimit = 5000

// changelogArgs normalises the optional cursor/limit args shared by every
// *Changelog query. after_id is the keyset tiebreaker (nil → the zero UUID, so
// the first page starts at the beginning of the since timestamp).
func changelogArgs(afterID *uuid.UUID, limit *int) (uuid.UUID, int32) {
after := uuid.Nil
if afterID != nil {
after = *afterID
}
l := defaultChangelogLimit
if limit != nil {
l = *limit
}
return after, int32(l)
}

func (r *queryResolver) SceneChangelog(ctx context.Context, since time.Time, afterID *uuid.UUID, limit *int) ([]models.EntityChange, error) {
after, l := changelogArgs(afterID, limit)
return r.services.Scene().Changelog(ctx, since, after, l)
}

func (r *queryResolver) PerformerChangelog(ctx context.Context, since time.Time, afterID *uuid.UUID, limit *int) ([]models.EntityChange, error) {
after, l := changelogArgs(afterID, limit)
return r.services.Performer().Changelog(ctx, since, after, l)
}

func (r *queryResolver) StudioChangelog(ctx context.Context, since time.Time, afterID *uuid.UUID, limit *int) ([]models.EntityChange, error) {
after, l := changelogArgs(afterID, limit)
return r.services.Studio().Changelog(ctx, since, after, l)
}

func (r *queryResolver) TagChangelog(ctx context.Context, since time.Time, afterID *uuid.UUID, limit *int) ([]models.EntityChange, error) {
after, l := changelogArgs(afterID, limit)
return r.services.Tag().Changelog(ctx, since, after, l)
}
8 changes: 8 additions & 0 deletions internal/converter/converter.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,14 @@ var (
updateParamsConverter = &gen.UpdateParamsConverterImpl{}
)

// NullUUIDToPtr converts a uuid.NullUUID to a *uuid.UUID, nil when invalid.
func NullUUIDToPtr(n uuid.NullUUID) *uuid.UUID {
if n.Valid {
return &n.UUID
}
return nil
}

// ImageToModel converts a queries.Image to a models.Image
func ImageToModel(i queries.Image) models.Image {
return modelConverter.ConvertImage(i)
Expand Down
2 changes: 1 addition & 1 deletion internal/database/database.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ import (

const (
postgresDriver = "postgres"
schemaVersion = 72
schemaVersion = 73
)

//go:embed migrations/postgres/*.sql
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
-- Non-partial keyset indexes for the entity changelog feed. Unlike the existing
-- partial (WHERE deleted = false) sort indexes, these must include deleted rows
-- so the changelog can report tombstones. Ordered (updated_at, id) ASC to match
-- the forward keyset scan; deleted is a payload column for index-only filtering.
CREATE INDEX scenes_updated_at_id_idx ON scenes (updated_at, id) INCLUDE (deleted);
CREATE INDEX performers_updated_at_id_idx ON performers (updated_at, id) INCLUDE (deleted);
CREATE INDEX studios_updated_at_id_idx ON studios (updated_at, id) INCLUDE (deleted);
CREATE INDEX tags_updated_at_id_idx ON tags (updated_at, id) INCLUDE (deleted);
Loading
Loading