diff --git a/db/attachment_test.go b/db/attachment_test.go index 16a3723c1a..17b083c7f1 100644 --- a/db/attachment_test.go +++ b/db/attachment_test.go @@ -1987,7 +1987,7 @@ func TestWritePathsPreserveUnmigratedAttachmentMetadata(t *testing.T) { return rev2ID }, run: func(t *testing.T, ctx context.Context, _ *Database, collection *DatabaseCollectionWithUser, docID, _ string) { - compacted, err := collection.CompactDocChannelHistory(ctx, docID, 100) + compacted, err := collection.CompactDocChannelHistory(ctx, docID, 100, []string{}) require.NoError(t, err) require.NotEmpty(t, compacted, "compaction should have pruned channel history, or no write happens") }, diff --git a/db/crud.go b/db/crud.go index 996580c383..76c8e4bd65 100644 --- a/db/crud.go +++ b/db/crud.go @@ -288,8 +288,9 @@ func (c *DatabaseCollection) GetDocChannelHistory(ctx context.Context, docid str } // CompactDocChannelHistory removes channel history entries that ended at or before the given sequence number. +// If seq is zero, it removes the history entries of the named channels instead. // This is used to prune stale channel assignment history to reduce storage overhead. -func (c *DatabaseCollection) CompactDocChannelHistory(ctx context.Context, docid string, seq uint64) ([]string, error) { +func (c *DatabaseCollection) CompactDocChannelHistory(ctx context.Context, docid string, seq uint64, channels []string) ([]string, error) { key := realDocID(docid) if key == "" { return nil, base.HTTPErrorf(400, "Invalid doc ID") @@ -326,7 +327,12 @@ func (c *DatabaseCollection) CompactDocChannelHistory(ctx context.Context, docid compactedChannels := make(base.Set) doc.SyncData.ChannelSetHistory = slices.DeleteFunc(doc.SyncData.ChannelSetHistory, func(channel ChannelSetEntry) bool { - del := channel.End <= seq + var del bool + if seq != 0 { + del = channel.End <= seq + } else if len(channels) > 0 { + del = slices.Contains(channels, channel.Name) + } if del { compactedChannels.Add(channel.Name) } @@ -334,7 +340,12 @@ func (c *DatabaseCollection) CompactDocChannelHistory(ctx context.Context, docid }) doc.SyncData.ChannelSet = slices.DeleteFunc(doc.SyncData.ChannelSet, func(channel ChannelSetEntry) bool { - del := channel.End != 0 && channel.End <= seq + var del bool + if seq != 0 { + del = channel.End != 0 && channel.End <= seq + } else if len(channels) > 0 { + del = channel.End != 0 && slices.Contains(channels, channel.Name) + } if del { compactedChannels.Add(channel.Name) } @@ -342,7 +353,16 @@ func (c *DatabaseCollection) CompactDocChannelHistory(ctx context.Context, docid }) for chanName, chanEntry := range doc.SyncData.Channels { - if chanEntry != nil && chanEntry.Seq <= seq { + if chanEntry == nil { + continue + } + var del bool + if seq != 0 { + del = chanEntry.Seq <= seq + } else if len(channels) > 0 { + del = slices.Contains(channels, chanName) + } + if del { compactedChannels.Add(chanName) delete(doc.SyncData.Channels, chanName) } diff --git a/db/hybrid_logical_vector_test.go b/db/hybrid_logical_vector_test.go index c2b045336c..9b6d3e8490 100644 --- a/db/hybrid_logical_vector_test.go +++ b/db/hybrid_logical_vector_test.go @@ -671,7 +671,7 @@ func TestRestampVersionCASMou(t *testing.T) { // A second metadata-only write, from a different path - channel history compaction, which unlike // resync applies to a tombstone as well as a live document. Both previous values have to be // carried forward, or they stop naming the last write to the body. - compacted, err := collection.CompactDocChannelHistory(ctx, docID, 100) + compacted, err := collection.CompactDocChannelHistory(ctx, docID, 100, []string{}) require.NoError(t, err) require.NotEmpty(t, compacted, "compaction should have pruned channel history, or no write happens") diff --git a/db/import_test.go b/db/import_test.go index a73fca5bea..db3d9d4525 100644 --- a/db/import_test.go +++ b/db/import_test.go @@ -1948,7 +1948,7 @@ func TestMetadataOnlyUpdateWritePaths(t *testing.T) { return rev2ID }, run: func(t *testing.T, ctx context.Context, collection *DatabaseCollectionWithUser, docID, _ string) { - compacted, err := collection.CompactDocChannelHistory(ctx, docID, 100) + compacted, err := collection.CompactDocChannelHistory(ctx, docID, 100, []string{""}) require.NoError(t, err) require.NotEmpty(t, compacted, "compaction should have pruned channel history, or no write happens") }, diff --git a/docs/api/paths/admin/keyspace-_channel_history-compact.yaml b/docs/api/paths/admin/keyspace-_channel_history-compact.yaml index 01efbbba51..44efab7c8a 100644 --- a/docs/api/paths/admin/keyspace-_channel_history-compact.yaml +++ b/docs/api/paths/admin/keyspace-_channel_history-compact.yaml @@ -11,10 +11,13 @@ parameters: post: summary: Compact Channel History of Document description: |- - Compacts channel history for a specified document. Channel history older than the specified sequence will be removed. + Compacts channel history for a specified document. The request must contain either `seq` or `channels`, but not both. - This endpoint removes all channel entries (for sequences before the specified sequence number where the document left the channel), - effectively cleaning up historical channel membership information while preserving active channels and recent changes. + With `seq`, this endpoint removes all channel entries where the document left the channel at or before the specified sequence number. + With `channels`, this endpoint removes the channel history entries for the named channels, whatever sequence the document left them at. + Channels that the document is still in are not removed. + + Compaction cleans up historical channel membership information while preserving active channels and recent changes. This can be useful for reducing metadata size for documents that frequently gain and lose access to channels. Required Sync Gateway RBAC roles: @@ -25,23 +28,42 @@ post: application/json: schema: type: object - required: - - seq properties: seq: description: |- - Channel history having end sequences earlier than this sequence will be removed from the specified document's metadata. + Channel history having end sequences at or before this sequence will be removed from the specified document's metadata. + + Cannot be used together with `channels`. type: integer format: int64 minimum: 1 example: 12345 + channels: + description: |- + Channel history for these channels will be removed from the specified document's metadata. + + A non-empty array cannot be used together with `seq`. + type: array + items: + type: string + example: + - channel1 + - channel2 + oneOf: + - required: + - seq + - required: + - channels + properties: + channels: + minItems: 1 responses: '200': description: |- Successfully compacted channel history from the specified document. Returns a list of channels that were compacted. - If the response has an empty array, it means either no channels were compacted. + If the response has an empty array, it means no channels were compacted. content: application/json: @@ -55,13 +77,13 @@ post: description: |- Array of channel names that were compacted. description: |- - A array of all the compacted channels + An array of all the compacted channels example: compacted_channels: - channel1 - channel2 '400': - description: 'Bad request. This could be due to invalid request parameters such as invalid seq value.' + description: 'Bad request. This could be due to invalid request parameters such as an invalid seq value, an empty channels array with no seq, or both seq and a non-empty channels array in the same request.' content: application/json: schema: diff --git a/rest/api_test.go b/rest/api_test.go index b9f53b2c4f..20c8f06940 100644 --- a/rest/api_test.go +++ b/rest/api_test.go @@ -3873,6 +3873,20 @@ func TestFetchBackupRevisionByCVThroughAPI(t *testing.T) { assert.Nil(t, body[db.BodyCV]) } +// GetCompactionRequestFromType builds a compact request for the given mode. The "_seq" mode sets only +// Seq, and any other mode sets only Channels, because a request cannot contain both. +func GetCompactionRequestFromType(compactionType string, seq uint64, channelsToCompact []string) CompactDocChannelHistoryRequest { + req := CompactDocChannelHistoryRequest{} + if compactionType == "_seq" { + req.Seq = seq + } else { + req.Channels = channelsToCompact + } + return req +} + +// TestDocumentChannelHistoryCompact verifies document channel history compaction, both with the +// collection method and with the REST endpoint. Each case runs once by sequence and once by channel name. func TestDocumentChannelHistoryCompact(t *testing.T) { defer db.SuspendSequenceBatching()() @@ -3881,245 +3895,300 @@ func TestDocumentChannelHistoryCompact(t *testing.T) { collection, ctx := rt.GetSingleTestDatabaseCollection() - t.Run("basic compaction", func(t *testing.T) { - // Create a document with a single channel assignment - version := rt.PutDoc("doc1", `{"channels": ["test"]}`) - syncData, err := collection.GetDocSyncData(ctx, "doc1") - assert.NoError(t, err) - - require.Len(t, syncData.ChannelSet, 1) - assert.Equal(t, db.ChannelSetEntry{Name: "test", Start: 1, End: 0}, syncData.ChannelSet[0]) - assert.Len(t, syncData.ChannelSetHistory, 0) - - // Remove all channels - ends the existing channel range in ChannelSet - version = rt.UpdateDoc("doc1", version, `{"channels": []}`) - syncData, err = collection.GetDocSyncData(ctx, "doc1") - assert.NoError(t, err) + for _, test := range []struct { + name string + }{ + {"_seq"}, + {"_channels"}, + } { + t.Run("basic compaction"+test.name, func(t *testing.T) { + var channelsToCompact []string + // Create a document with a single channel assignment + docID := "basic_compaction" + test.name + version := rt.PutDoc(docID, `{"channels": ["test"]}`) + docAddedSeq := rt.GetDocumentSequence(docID) + syncData, err := collection.GetDocSyncData(ctx, docID) + assert.NoError(t, err) + + require.Len(t, syncData.ChannelSet, 1) + assert.Equal(t, db.ChannelSetEntry{Name: "test", Start: docAddedSeq, End: 0}, syncData.ChannelSet[0]) + assert.Len(t, syncData.ChannelSetHistory, 0) + + // Remove all channels - ends the existing channel range in ChannelSet + version = rt.UpdateDoc(docID, version, `{"channels": []}`) + syncData, err = collection.GetDocSyncData(ctx, docID) + compactionSeq := rt.GetDocumentSequence(docID) + channelsToCompact = append(channelsToCompact, "test") + assert.NoError(t, err) + + require.Len(t, syncData.ChannelSet, 1) + assert.Equal(t, db.ChannelSetEntry{Name: "test", Start: docAddedSeq, End: compactionSeq}, syncData.ChannelSet[0]) + assert.Len(t, syncData.ChannelSetHistory, 0) + + // Add multiple channels - generates history for the previously removed channel + _ = rt.UpdateDoc(docID, version, `{"channels": ["test", "test2"]}`) + syncData, err = collection.GetDocSyncData(ctx, docID) + assert.NoError(t, err) + channelsAddedAtSeq := rt.GetDocumentSequence(docID) + + require.Len(t, syncData.ChannelSet, 2) + assert.Contains(t, syncData.ChannelSet, db.ChannelSetEntry{Name: "test", Start: channelsAddedAtSeq, End: 0}) + assert.Contains(t, syncData.ChannelSet, db.ChannelSetEntry{Name: "test2", Start: channelsAddedAtSeq, End: 0}) + require.Len(t, syncData.ChannelSetHistory, 1) + assert.Equal(t, db.ChannelSetEntry{Name: "test", Start: docAddedSeq, End: compactionSeq}, syncData.ChannelSetHistory[0]) + + req := GetCompactionRequestFromType(test.name, compactionSeq, channelsToCompact) + + // Compact the removed channel's history and verify it is cleared + compactedChannels, err := collection.CompactDocChannelHistory(ctx, docID, req.Seq, req.Channels) + require.NoError(t, err) + assert.Equal(t, []string{"test"}, compactedChannels) + syncData, err = collection.GetDocSyncData(ctx, docID) + assert.NoError(t, err) + assert.Zero(t, len(syncData.ChannelSetHistory)) + }) - require.Len(t, syncData.ChannelSet, 1) - assert.Equal(t, db.ChannelSetEntry{Name: "test", Start: 1, End: 2}, syncData.ChannelSet[0]) - assert.Len(t, syncData.ChannelSetHistory, 0) + t.Run("seq zero and empty channels array keeps all history"+test.name, func(t *testing.T) { + docID := "no_compaction" + test.name + // An empty request (seq 0, no channels) keeps all entries + version := rt.PutDoc(docID, `{"channels": ["test"]}`) + version = rt.UpdateDoc(docID, version, `{"channels": []}`) + _ = rt.UpdateDoc(docID, version, `{"channels": ["test", "test2"]}`) - // Add multiple channels - generates history for the previously removed channel - _ = rt.UpdateDoc("doc1", version, `{"channels": ["test", "test2"]}`) - syncData, err = collection.GetDocSyncData(ctx, "doc1") - assert.NoError(t, err) + syncDataBefore, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) + historyLenBefore := len(syncDataBefore.ChannelSetHistory) + require.Greater(t, historyLenBefore, 0) - require.Len(t, syncData.ChannelSet, 2) - assert.Contains(t, syncData.ChannelSet, db.ChannelSetEntry{Name: "test", Start: 3, End: 0}) - assert.Contains(t, syncData.ChannelSet, db.ChannelSetEntry{Name: "test2", Start: 3, End: 0}) - require.Len(t, syncData.ChannelSetHistory, 1) - assert.Equal(t, db.ChannelSetEntry{Name: "test", Start: 1, End: 2}, syncData.ChannelSetHistory[0]) + req := GetCompactionRequestFromType(test.name, 0, []string{}) + compactedChannels, err := collection.CompactDocChannelHistory(ctx, docID, req.Seq, req.Channels) + require.NoError(t, err) + assert.Equal(t, []string{}, compactedChannels) - // Compact history at sequence 2 and verify history is cleared - compactedChannels, err := collection.CompactDocChannelHistory(ctx, "doc1", 2) - require.NoError(t, err) - assert.Equal(t, []string{"test"}, compactedChannels) - syncData, err = collection.GetDocSyncData(ctx, "doc1") - assert.NoError(t, err) - assert.Zero(t, len(syncData.ChannelSetHistory)) - }) + syncDataAfter, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) + assert.Equal(t, historyLenBefore, len(syncDataAfter.ChannelSetHistory)) + }) - t.Run("seq zero keeps all history", func(t *testing.T) { - // Compact with seq=0 should keep all entries - version := rt.PutDoc("doc2", `{"channels": ["test"]}`) - version = rt.UpdateDoc("doc2", version, `{"channels": []}`) - _ = rt.UpdateDoc("doc2", version, `{"channels": ["test", "test2"]}`) + t.Run("compact all history"+test.name, func(t *testing.T) { + var channelsToCompact []string + // A very high seq, or the channel name, removes all history + docID := "compact_all_history" + test.name + version := rt.PutDoc(docID, `{"channels": ["test"]}`) + channelsToCompact = append(channelsToCompact, "test") + version = rt.UpdateDoc(docID, version, `{"channels": []}`) + _ = rt.UpdateDoc(docID, version, `{"channels": ["test", "test2"]}`) - syncDataBefore, err := collection.GetDocSyncData(ctx, "doc2") - require.NoError(t, err) - historyLenBefore := len(syncDataBefore.ChannelSetHistory) - require.Greater(t, historyLenBefore, 0) + syncDataBefore, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) - compactedChannels, err := collection.CompactDocChannelHistory(ctx, "doc2", 0) - require.NoError(t, err) - assert.Equal(t, []string{}, compactedChannels) + req := GetCompactionRequestFromType(test.name, 9999999, channelsToCompact) + compactedChannels, err := collection.CompactDocChannelHistory(ctx, docID, req.Seq, req.Channels) + require.NoError(t, err) + assert.Equal(t, []string{"test"}, compactedChannels) - syncDataAfter, err := collection.GetDocSyncData(ctx, "doc2") - require.NoError(t, err) - assert.Equal(t, historyLenBefore, len(syncDataAfter.ChannelSetHistory)) - }) + syncData, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) + assert.Zero(t, len(syncData.ChannelSetHistory)) + assert.Equal(t, len(syncDataBefore.Channels), len(syncData.Channels)) + assert.Equal(t, len(syncDataBefore.ChannelSet), len(syncData.ChannelSet)) + }) - t.Run("compact all history", func(t *testing.T) { - // Compact with very high seq removes all history - version := rt.PutDoc("doc3", `{"channels": ["test"]}`) - version = rt.UpdateDoc("doc3", version, `{"channels": []}`) - _ = rt.UpdateDoc("doc3", version, `{"channels": ["test", "test2"]}`) + t.Run("partial history compaction"+test.name, func(t *testing.T) { + docID := "partial_history_compaction" + test.name + var channelsToCompact []string + // Compact removes entries with End <= seq (or the named channels), keeps the rest + version := rt.PutDoc(docID, `{"channels": ["a"]}`) + version = rt.UpdateDoc(docID, version, `{"channels": []}`) - syncDataBefore, err := collection.GetDocSyncData(ctx, "doc3") - require.NoError(t, err) + docSeq := rt.GetDocumentSequence(docID) + channelsToCompact = append(channelsToCompact, "a") - // Compact with very high seq number - compactedChannels, err := collection.CompactDocChannelHistory(ctx, "doc3", 999999) - require.NoError(t, err) - assert.Equal(t, []string{"test"}, compactedChannels) + version = rt.UpdateDoc(docID, version, `{"channels": ["b"]}`) + _ = rt.UpdateDoc(docID, version, `{"channels": []}`) - syncData, err := collection.GetDocSyncData(ctx, "doc3") - require.NoError(t, err) - assert.Zero(t, len(syncData.ChannelSetHistory)) - assert.Equal(t, len(syncDataBefore.Channels), len(syncData.Channels)) - assert.Equal(t, len(syncDataBefore.ChannelSet), len(syncData.ChannelSet)) - }) + syncDataBefore, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) + channelSetLenBefore := len(syncDataBefore.ChannelSet) + channelLenBefore := len(syncDataBefore.Channels) - t.Run("partial history compaction", func(t *testing.T) { - // Compact removes entries with End <= seq, keeps entries with End > seq - version := rt.PutDoc("doc4", `{"channels": ["a"]}`) - version = rt.UpdateDoc("doc4", version, `{"channels": []}`) + req := GetCompactionRequestFromType(test.name, docSeq, channelsToCompact) + compactedChannels, err := collection.CompactDocChannelHistory(ctx, docID, req.Seq, req.Channels) + require.NoError(t, err) + assert.Equal(t, []string{"a"}, compactedChannels) - docSeq := rt.GetDocumentSequence("doc4") + syncDataAfter, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) - version = rt.UpdateDoc("doc4", version, `{"channels": ["b"]}`) - _ = rt.UpdateDoc("doc4", version, `{"channels": []}`) + for _, entry := range syncDataAfter.ChannelSet { + assert.True(t, entry.End > docSeq) + } - syncDataBefore, err := collection.GetDocSyncData(ctx, "doc4") - require.NoError(t, err) - channelSetLenBefore := len(syncDataBefore.ChannelSet) - channelLenBefore := len(syncDataBefore.Channels) + for _, entry := range syncDataAfter.Channels { + assert.True(t, entry.Seq > docSeq) + } - compactedChannels, err := collection.CompactDocChannelHistory(ctx, "doc4", docSeq) - require.NoError(t, err) - assert.Equal(t, []string{"a"}, compactedChannels) + // Should have fewer than before + assert.Less(t, len(syncDataAfter.Channels), channelLenBefore) + assert.Less(t, len(syncDataAfter.ChannelSet), channelSetLenBefore) + }) - syncDataAfter, err := collection.GetDocSyncData(ctx, "doc4") - require.NoError(t, err) + t.Run("nonexistent document"+test.name, func(t *testing.T) { + docID := "nonexistent" + test.name + req := GetCompactionRequestFromType(test.name, 1, []string{"test"}) + // Compacting nonexistent doc should return not found error + _, err := collection.CompactDocChannelHistory(ctx, docID, req.Seq, req.Channels) + assert.Error(t, err) - for _, entry := range syncDataAfter.ChannelSet { - assert.True(t, entry.End == 0 || entry.End > docSeq) - } + bodyBytes, err := base.JSONMarshal(req) + require.NoError(t, err) + resp := rt.SendAdminRequest("POST", "/{{.keyspace}}/_channel_history/nonexistent/compact", string(bodyBytes)) + RequireStatus(t, resp, http.StatusNotFound) + }) - for _, entry := range syncDataAfter.Channels { - assert.True(t, entry.Seq > docSeq) - } + t.Run("multiple channels with mixed history"+test.name, func(t *testing.T) { + docID := "multiple_channels_with_mixed_history" + test.name + // Complex scenario with multiple channels + version := rt.PutDoc(docID, `{"channels": ["a", "b"]}`) + version = rt.UpdateDoc(docID, version, `{"channels": ["a"]}`) - // Should have fewer than before - assert.Less(t, len(syncDataAfter.Channels), channelLenBefore) - assert.Less(t, len(syncDataAfter.ChannelSet), channelSetLenBefore) - }) + docSeq := rt.GetDocumentSequence(docID) + channelsToCompact := []string{"b"} - t.Run("invalid doc id", func(t *testing.T) { - // Empty doc ID should return 400 error - _, err := collection.CompactDocChannelHistory(ctx, "", 1) - assert.Error(t, err) - assert.Contains(t, err.Error(), "Invalid doc ID") - }) + _ = rt.UpdateDoc(docID, version, `{"channels": ["a", "c"]}`) - t.Run("nonexistent document", func(t *testing.T) { - // Compacting nonexistent doc should return not found error - _, err := collection.CompactDocChannelHistory(ctx, "nonexistent", 1) - assert.Error(t, err) + syncDataBefore, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) - req := CompactDocChannelHistoryRequest{ - Seq: 999999, - } + req := GetCompactionRequestFromType(test.name, docSeq, channelsToCompact) + compactedChannels, err := collection.CompactDocChannelHistory(ctx, docID, req.Seq, req.Channels) + require.NoError(t, err) + assert.Equal(t, []string{"b"}, compactedChannels) - bodyBytes, err := base.JSONMarshal(req) - require.NoError(t, err) - resp := rt.SendAdminRequest("POST", "/{{.keyspace}}/_channel_history/nonexistent/compact", string(bodyBytes)) - RequireStatus(t, resp, http.StatusNotFound) - }) + syncData, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) + // Entries for the compacted channel are removed + assert.Less(t, len(syncData.ChannelSet), len(syncDataBefore.ChannelSet)) + assert.Less(t, len(syncData.Channels), len(syncDataBefore.Channels)) + }) - t.Run("multiple channels with mixed history", func(t *testing.T) { - // Complex scenario with multiple channels - version := rt.PutDoc("doc7", `{"channels": ["a", "b"]}`) - version = rt.UpdateDoc("doc7", version, `{"channels": ["a"]}`) + t.Run("compact empty history"+test.name, func(t *testing.T) { + docID := "compact_empty_history" + test.name + // Doc with no history should succeed + rt.PutDoc(docID, `{"channels": ["test"]}`) - docSeq := rt.GetDocumentSequence("doc7") + syncDataBefore, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) + assert.Len(t, syncDataBefore.ChannelSetHistory, 0) - _ = rt.UpdateDoc("doc7", version, `{"channels": ["a", "c"]}`) + req := GetCompactionRequestFromType(test.name, 1, []string{"test"}) - syncDataBefore, err := collection.GetDocSyncData(ctx, "doc7") - require.NoError(t, err) - compactedChannels, err := collection.CompactDocChannelHistory(ctx, "doc7", docSeq) - require.NoError(t, err) - assert.Equal(t, []string{"b"}, compactedChannels) + compactedChannels, err := collection.CompactDocChannelHistory(ctx, docID, req.Seq, req.Channels) + require.NoError(t, err) + assert.Equal(t, []string{}, compactedChannels) - syncData, err := collection.GetDocSyncData(ctx, "doc7") - require.NoError(t, err) - // Should still have active channels - assert.Less(t, len(syncData.ChannelSet), len(syncDataBefore.ChannelSet)) - assert.Less(t, len(syncData.Channels), len(syncDataBefore.Channels)) - }) + syncDataAfter, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) + assert.Len(t, syncDataAfter.ChannelSetHistory, 0) + }) - t.Run("compact empty history", func(t *testing.T) { - // Doc with no history should succeed - rt.PutDoc("doc8", `{"channels": ["test"]}`) + t.Run("verify xattr updates"+test.name, func(t *testing.T) { + docID := "verify_xattr_updates" + test.name + version := rt.PutDoc(docID, `{"channels": ["test"]}`) + version = rt.UpdateDoc(docID, version, `{"channels": []}`) + _ = rt.UpdateDoc(docID, version, `{"channels": ["test", "test2"]}`) + seq := rt.GetDocumentSequence(docID) + channelsToCompact := []string{"test"} - syncDataBefore, err := collection.GetDocSyncData(ctx, "doc8") - require.NoError(t, err) - assert.Len(t, syncDataBefore.ChannelSetHistory, 0) + req := GetCompactionRequestFromType(test.name, seq, channelsToCompact) - compactedChannels, err := collection.CompactDocChannelHistory(ctx, "doc8", 1) - require.NoError(t, err) - assert.Equal(t, []string{}, compactedChannels) + // After compaction, document should still be accessible + compactedChannels, err := collection.CompactDocChannelHistory(ctx, docID, req.Seq, req.Channels) + require.NoError(t, err) + assert.Equal(t, []string{"test"}, compactedChannels) - syncDataAfter, err := collection.GetDocSyncData(ctx, "doc8") - require.NoError(t, err) - assert.Len(t, syncDataAfter.ChannelSetHistory, 0) - }) + // Verify doc is still accessible with correct data + syncData, err := collection.GetDocSyncData(ctx, docID) + require.NoError(t, err) + assert.NotNil(t, syncData) + }) - t.Run("verify xattr updates", func(t *testing.T) { - // Verify CAS is updated and no reimport happens - version := rt.PutDoc("doc10", `{"channels": ["test"]}`) - version = rt.UpdateDoc("doc10", version, `{"channels": []}`) - _ = rt.UpdateDoc("doc10", version, `{"channels": ["test", "test2"]}`) - seq := rt.GetDocumentSequence("doc10") + t.Run("test rest api endpoint"+test.name, func(t *testing.T) { + docID1 := "rest_api_endpoint_1" + test.name + docID2 := "rest_api_endpoint_2" + test.name + version := rt.PutDoc(docID1, `{"channels": ["test"]}`) + version = rt.UpdateDoc(docID1, version, `{"channels": []}`) + _ = rt.UpdateDoc(docID1, version, `{"channels": ["test", "test2"]}`) - // After compaction, document should still be accessible - compactedChannels, err := collection.CompactDocChannelHistory(ctx, "doc10", seq) - require.NoError(t, err) - assert.Equal(t, []string{"test"}, compactedChannels) + version = rt.PutDoc(docID2, `{"channels": ["a", "b"]}`) + version = rt.UpdateDoc(docID2, version, `{"channels": ["a"]}`) + _ = rt.UpdateDoc(docID2, version, `{"channels": ["a", "c"]}`) - // Verify doc is still accessible with correct data - syncData, err := collection.GetDocSyncData(ctx, "doc10") - require.NoError(t, err) - assert.NotNil(t, syncData) - }) + req1 := GetCompactionRequestFromType(test.name, 999999, []string{"test"}) + bodyBytes1, err := base.JSONMarshal(req1) + require.NoError(t, err) + resp := rt.SendAdminRequest("POST", fmt.Sprintf("/{{.keyspace}}/_channel_history/%s/compact", docID1), string(bodyBytes1)) + RequireStatus(t, resp, http.StatusOK) - t.Run("test rest api endpoint", func(t *testing.T) { + var chanOutput1 map[string][]string + err = base.JSONUnmarshal(resp.Body.Bytes(), &chanOutput1) + require.NoError(t, err) - version := rt.PutDoc("doc11", `{"channels": ["test"]}`) - version = rt.UpdateDoc("doc11", version, `{"channels": []}`) - _ = rt.UpdateDoc("doc11", version, `{"channels": ["test", "test2"]}`) + expectedOutput1 := map[string][]string{ + "compacted_channels": []string{"test"}, + } - version = rt.PutDoc("doc12", `{"channels": ["a", "b"]}`) - version = rt.UpdateDoc("doc12", version, `{"channels": ["a"]}`) - _ = rt.UpdateDoc("doc12", version, `{"channels": ["a", "c"]}`) - req := CompactDocChannelHistoryRequest{ - Seq: 999999, - } + assert.Equal(t, expectedOutput1, chanOutput1) - bodyBytes, err := base.JSONMarshal(req) - require.NoError(t, err) - resp := rt.SendAdminRequest("POST", "/{{.keyspace}}/_channel_history/doc11/compact", string(bodyBytes)) - RequireStatus(t, resp, http.StatusOK) + req2 := GetCompactionRequestFromType(test.name, 999999, []string{"b"}) + bodyBytes2, err := base.JSONMarshal(req2) + require.NoError(t, err) + resp2 := rt.SendAdminRequest("POST", fmt.Sprintf("/{{.keyspace}}/_channel_history/%s/compact", docID2), string(bodyBytes2)) + RequireStatus(t, resp2, http.StatusOK) - var chanOutput1 map[string][]string - err = base.JSONUnmarshal(resp.Body.Bytes(), &chanOutput1) - require.NoError(t, err) + var chanOutput2 map[string][]string + err = base.JSONUnmarshal(resp2.Body.Bytes(), &chanOutput2) + require.NoError(t, err) - expectedOutput1 := map[string][]string{ - "compacted_channels": []string{"test"}, - } + expectedOutput2 := map[string][]string{ + "compacted_channels": []string{"b"}, + } - assert.Equal(t, expectedOutput1, chanOutput1) + assert.Equal(t, expectedOutput2, chanOutput2) + }) + } - resp2 := rt.SendAdminRequest("POST", "/{{.keyspace}}/_channel_history/doc12/compact", string(bodyBytes)) - RequireStatus(t, resp2, http.StatusOK) + t.Run("invalid doc id", func(t *testing.T) { + // Empty doc ID should return 400 error + _, err := collection.CompactDocChannelHistory(ctx, "", 1, []string{}) + assert.Error(t, err) + assert.Contains(t, err.Error(), "Invalid doc ID") + }) - var chanOutput2 map[string][]string - err = base.JSONUnmarshal(resp2.Body.Bytes(), &chanOutput2) + t.Run("empty request returns 400", func(t *testing.T) { + req := CompactDocChannelHistoryRequest{ + Seq: 0, + } + bodyBytes, err := base.JSONMarshal(req) require.NoError(t, err) + resp := rt.SendAdminRequest("POST", "/{{.keyspace}}/_channel_history/doc1/compact", string(bodyBytes)) + RequireStatus(t, resp, http.StatusBadRequest) + }) - expectedOutput2 := map[string][]string{ - "compacted_channels": []string{"b"}, + t.Run("empty channels returns 400", func(t *testing.T) { + req := CompactDocChannelHistoryRequest{ + Channels: []string{}, } - - assert.Equal(t, expectedOutput2, chanOutput2) + bodyBytes, err := base.JSONMarshal(req) + require.NoError(t, err) + resp := rt.SendAdminRequest("POST", "/{{.keyspace}}/_channel_history/doc1/compact", string(bodyBytes)) + RequireStatus(t, resp, http.StatusBadRequest) }) - t.Run("seq zero returns 400", func(t *testing.T) { + t.Run("passing both sequence and channels returns 400", func(t *testing.T) { req := CompactDocChannelHistoryRequest{ - Seq: 0, + Seq: 1, + Channels: []string{"chanA"}, } bodyBytes, err := base.JSONMarshal(req) require.NoError(t, err) @@ -4135,9 +4204,10 @@ func TestCompactNonImportedDocWithAutoImport(t *testing.T) { // Create RestTester with AutoImport disabled to allow non-imported documents rtConfig := RestTesterConfig{ SyncFn: channels.DocChannelsSyncFunction, - DatabaseConfig: &DatabaseConfig{DbConfig: DbConfig{ - AutoImport: false, - }}, + DatabaseConfig: &DatabaseConfig{ + DbConfig: DbConfig{ + AutoImport: false, + }}, } rt := NewRestTesterDefaultCollection(t, &rtConfig) defer rt.Close() @@ -4145,84 +4215,94 @@ func TestCompactNonImportedDocWithAutoImport(t *testing.T) { collection, ctx := rt.GetSingleTestDatabaseCollection() dataStore := rt.GetSingleDataStore() - // Step 1-6: Use REST API to create document with channel history - nonImportedDocID := "non_imported_doc" + for _, test := range []struct { + name string + }{ + {"_seq"}, + {"_channels"}, + } { + t.Run(test.name, func(t *testing.T) { - // Create initial document with a channel - version := rt.PutDoc(nonImportedDocID, `{"type":"test","channels":["test_channel"]}`) + // Use the REST API to create a document with channel history + nonImportedDocID := "non_imported_doc" + test.name - // Update to remove channel (creates history) - version = rt.UpdateDoc(nonImportedDocID, version, `{"type":"test","channels":[]}`) + // Create initial document with a channel + version := rt.PutDoc(nonImportedDocID, `{"type":"test","channels":["test_channel"]}`) - // Update again to add new channels (more history) - _ = rt.UpdateDoc(nonImportedDocID, version, `{"type":"test","channels":["test_channel","new_channel"]}`) + // Update to remove channel (creates history) + version = rt.UpdateDoc(nonImportedDocID, version, `{"type":"test","channels":[]}`) - // Verify document has channel history - syncDataBefore, err := collection.GetDocSyncData(ctx, nonImportedDocID) - require.NoError(t, err) - require.Greater(t, len(syncDataBefore.ChannelSetHistory), 0, "document should have channel history") - cvBeforeCompaction := syncDataBefore.CVOrRevTreeID() - require.NotEmpty(t, cvBeforeCompaction, "cv should be set before compaction") + // Update again to add new channels (more history) + _ = rt.UpdateDoc(nonImportedDocID, version, `{"type":"test","channels":["test_channel","new_channel"]}`) - // Step 6: Get document sequence for compaction point - docSeq := rt.GetDocumentSequence(nonImportedDocID) + // Verify document has channel history + syncDataBefore, err := collection.GetDocSyncData(ctx, nonImportedDocID) + require.NoError(t, err) + require.Greater(t, len(syncDataBefore.ChannelSetHistory), 0, "document should have channel history") + cvBeforeCompaction := syncDataBefore.CVOrRevTreeID() + require.NotEmpty(t, cvBeforeCompaction, "cv should be set before compaction") - // Step 7: Update the document body directly in the datastore to simulate external modification - // Read current document body - docBytesRaw, _, err := dataStore.GetRaw(ctx, nonImportedDocID) - require.NoError(t, err) - var docBody map[string]any - err = json.Unmarshal(docBytesRaw, &docBody) - require.NoError(t, err) + // Get the document sequence for the compaction point + docSeq := rt.GetDocumentSequence(nonImportedDocID) + channelsToCompact := []string{"test_channel"} - // Modify document body with external changes - docBody["modified"] = true - docBody["updatedAt"] = "external_update" - docBody["externalVersion"] = 2 - modifiedDocBytes, err := json.Marshal(docBody) - require.NoError(t, err) + // Update the document body directly in the datastore to simulate external modification + // Read current document body + docBytesRaw, _, err := dataStore.GetRaw(ctx, nonImportedDocID) + require.NoError(t, err) + var docBody map[string]any + err = json.Unmarshal(docBytesRaw, &docBody) + require.NoError(t, err) - // Write modified body back to datastore - err = dataStore.SetRaw(ctx, nonImportedDocID, 0, nil, modifiedDocBytes) - require.NoError(t, err) + // Modify document body with external changes + docBody["modified"] = true + docBody["updatedAt"] = "external_update" + docBody["externalVersion"] = 2 + modifiedDocBytes, err := json.Marshal(docBody) + require.NoError(t, err) - // Step 8: Call CompactDocChannelHistory - this will trigger the auto-import check - // which verifies the document is imported (has valid _sync xattr) before compacting - compactedChannels, err := collection.CompactDocChannelHistory(ctx, nonImportedDocID, docSeq-1) - require.NoError(t, err) - assert.Equal(t, []string{"test_channel"}, compactedChannels) + // Write modified body back to datastore + err = dataStore.SetRaw(ctx, nonImportedDocID, 0, nil, modifiedDocBytes) + require.NoError(t, err) - // Step 9: Verify compaction succeeded and history was removed - require.NoError(t, err) + // Call CompactDocChannelHistory - this will trigger the auto-import check + // which verifies the document is imported (has valid _sync xattr) before compacting + req := GetCompactionRequestFromType(test.name, docSeq-1, channelsToCompact) + compactedChannels, err := collection.CompactDocChannelHistory(ctx, nonImportedDocID, req.Seq, req.Channels) + require.NoError(t, err) + assert.Equal(t, []string{"test_channel"}, compactedChannels) - xattrs, cas, err := dataStore.GetXattrs(ctx, nonImportedDocID, []string{base.SyncXattrName, base.VvXattrName, base.MouXattrName}) - require.NoError(t, err) - doc := db.NewDocument(nonImportedDocID) - err = doc.UnmarshalWithXattrs(ctx, nil, xattrs[base.SyncXattrName], xattrs[base.VvXattrName], nil, nil, db.DocUnmarshalSync) - require.NoError(t, err) - err = base.JSONUnmarshal(xattrs[base.MouXattrName], &doc.MetadataOnlyUpdate) - require.NoError(t, err) + xattrs, cas, err := dataStore.GetXattrs(ctx, nonImportedDocID, []string{base.SyncXattrName, base.VvXattrName, base.MouXattrName}) + require.NoError(t, err) + doc := db.NewDocument(nonImportedDocID) + err = doc.UnmarshalWithXattrs(ctx, nil, xattrs[base.SyncXattrName], xattrs[base.VvXattrName], nil, nil, db.DocUnmarshalSync) + require.NoError(t, err) + err = base.JSONUnmarshal(xattrs[base.MouXattrName], &doc.MetadataOnlyUpdate) + require.NoError(t, err) - // History should be compacted - assert.Less(t, len(doc.SyncData.ChannelSetHistory), len(syncDataBefore.ChannelSetHistory)) - // CV must be updated by the import triggered during compaction - assert.NotEqual(t, cvBeforeCompaction, doc.SyncData.CVOrRevTreeID(), "cv should be updated after compaction import") - // sync cas should be equal to doc cas - assert.Equal(t, doc.SyncData.Cas, base.CasToString(cas)) - // verify _mou.cas - mouCAS := base.HexCasToUint64(doc.MetadataOnlyUpdate.HexCAS) - assert.Equal(t, mouCAS, cas) - - // Step 10: Verify document is still accessible and intact - docFromBucket, _, err := rt.GetSingleDataStore().GetRaw(ctx, nonImportedDocID) - require.NoError(t, err) - require.NotNil(t, docFromBucket) + // History should be compacted + assert.Less(t, len(doc.SyncData.ChannelSetHistory), len(syncDataBefore.ChannelSetHistory)) + // CV must be updated by the import triggered during compaction + assert.NotEqual(t, cvBeforeCompaction, doc.SyncData.CVOrRevTreeID(), "cv should be updated after compaction import") + // sync cas should be equal to doc cas + assert.Equal(t, doc.SyncData.Cas, base.CasToString(cas)) + // verify _mou.cas + mouCAS := base.HexCasToUint64(doc.MetadataOnlyUpdate.HexCAS) + assert.Equal(t, mouCAS, cas) + + // Verify the document is still accessible and intact + docFromBucket, _, err := rt.GetSingleDataStore().GetRaw(ctx, nonImportedDocID) + require.NoError(t, err) + require.NotNil(t, docFromBucket) - // Verify the document body is intact after compaction - var finalBody map[string]any - err = json.Unmarshal(docFromBucket, &finalBody) - require.NoError(t, err) - assert.Equal(t, "test", finalBody["type"]) + // Verify the document body is intact after compaction + var finalBody map[string]any + err = json.Unmarshal(docFromBucket, &finalBody) + require.NoError(t, err) + assert.Equal(t, "test", finalBody["type"]) + }) + + } } // TestGetDocChannelHistory tests the GetDocChannelHistory function and the diff --git a/rest/doc_api.go b/rest/doc_api.go index 427f7fc2df..c8681f45c4 100644 --- a/rest/doc_api.go +++ b/rest/doc_api.go @@ -959,13 +959,23 @@ func (h *handler) handleGetDocChannelHistory() error { } type CompactDocChannelHistoryRequest struct { - Seq uint64 `json:"seq"` + Seq uint64 `json:"seq"` + Channels []string `json:"channels"` +} + +func (r *CompactDocChannelHistoryRequest) validate() error { + if r.Seq != 0 && len(r.Channels) != 0 { + return base.HTTPErrorf(http.StatusBadRequest, "channels and seq cannot be used together") + } else if r.Seq == 0 && len(r.Channels) == 0 { + return base.HTTPErrorf(http.StatusBadRequest, "must specify channels or sequence") + } + return nil } // handleCompactDocChannelHistory handles POST /{keyspace}/_channel_history/{docid}/compact. -// It accepts a JSON body containing a sequence number and removes channel history entries -// for the specified document that ended at or before that sequence, returning the compacted -// channel names for that document. +// It accepts a JSON body containing either a sequence number or a list of channel names, but not both. +// It removes channel history entries for the specified document that ended at or before that sequence, +// or that belong to the named channels, and returns the compacted channel names. func (h *handler) handleCompactDocChannelHistory() error { h.assertAdminOnly() @@ -978,11 +988,12 @@ func (h *handler) handleCompactDocChannelHistory() error { return base.HTTPErrorf(http.StatusBadRequest, "invalid JSON: %v", err) } - if req.Seq == 0 { - return base.HTTPErrorf(http.StatusBadRequest, "missing seq") + err = req.validate() + if err != nil { + return err } - channels, err := h.collection.CompactDocChannelHistory(h.ctx(), docid, req.Seq) + channels, err := h.collection.CompactDocChannelHistory(h.ctx(), docid, req.Seq, req.Channels) if err != nil { return err } diff --git a/rest/doc_api_test.go b/rest/doc_api_test.go index 12cb927203..56c7fc33de 100644 --- a/rest/doc_api_test.go +++ b/rest/doc_api_test.go @@ -362,3 +362,42 @@ func TestCVUnescapedRevQueryParam(t *testing.T) { }) } } + +// TestCompactChannelHistoryRequestValidate verifies that a compact request must set exactly one of +// seq or channels, and that any other combination returns a 400 error. +func TestCompactChannelHistoryRequestValidate(t *testing.T) { + tests := []struct { + req CompactDocChannelHistoryRequest + err error + }{ + { + req: CompactDocChannelHistoryRequest{ + Seq: 0, + Channels: []string{}, + }, + err: base.HTTPErrorf(http.StatusBadRequest, "must specify channels or sequence"), + }, + { + req: CompactDocChannelHistoryRequest{ + Seq: 10, + Channels: []string{"chanA"}, + }, + err: base.HTTPErrorf(http.StatusBadRequest, "channels and seq cannot be used together"), + }, + { + req: CompactDocChannelHistoryRequest{ + Seq: 10, + }, + err: nil, + }, + { + req: CompactDocChannelHistoryRequest{ + Channels: []string{"chanA"}, + }, + err: nil, + }, + } + for _, test := range tests { + assert.Equal(t, test.err, test.req.validate()) + } +}