From 6b8b21745301cf6f0e6874cbfac8c4bdd607fdd0 Mon Sep 17 00:00:00 2001 From: Rob Rudin Date: Thu, 17 Sep 2026 16:21:27 -0400 Subject: [PATCH] MLE-32904 Incremental write now supportings matching on a field Supports a use case where the document to be matched against has had its URI modified since it was ingested. --- marklogic-client-api/build.gradle | 14 +++ .../filter/IncrementalWriteConfig.java | 23 +++- .../filter/IncrementalWriteFilter.java | 56 ++++++++- .../IncrementalWriteFromLexiconsFilter.java | 14 ++- .../IncrementalWriteFromViewFilter.java | 15 ++- .../filter/IncrementalWriteFilterTest.java | 34 ++++++ .../IncrementalWriteSourceUriKeyNameTest.java | 107 ++++++++++++++++++ .../filter/IncrementalWriteTest.java | 45 ++++++++ .../ml-config/databases/content-database.json | 16 +++ 9 files changed, 312 insertions(+), 12 deletions(-) create mode 100644 marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteSourceUriKeyNameTest.java diff --git a/marklogic-client-api/build.gradle b/marklogic-client-api/build.gradle index 58597c7de..59e2d8c81 100644 --- a/marklogic-client-api/build.gradle +++ b/marklogic-client-api/build.gradle @@ -3,6 +3,7 @@ */ plugins { + id 'net.saliman.properties' version '1.6.0' id 'maven-publish' } @@ -186,6 +187,19 @@ publishing { } } repositories { + // Supports publishing to a GH packages repository for GH Action workflows that are not able to access our + // internal repository. + if (project.hasProperty("ghToken") && project.hasProperty("ghPackagesUrl")) { + maven { + name = "GitHubPackages" + url = uri(ghPackagesUrl) + credentials { + username = ghActor + password = ghToken + } + } + } + maven { if (project.hasProperty("mavenUser")) { credentials { diff --git a/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteConfig.java b/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteConfig.java index d350f71f1..237740ed2 100644 --- a/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteConfig.java +++ b/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteConfig.java @@ -17,6 +17,7 @@ public class IncrementalWriteConfig { private final String hashKeyName; + private final String sourceUriKeyName; private final String timestampKeyName; private final boolean canonicalizeJson; private final Consumer skippedDocumentsConsumer; @@ -26,11 +27,24 @@ public class IncrementalWriteConfig { private final String schemaName; private final String viewName; - public IncrementalWriteConfig(String hashKeyName, String timestampKeyName, boolean canonicalizeJson, + // This was mistakenly exposed as a public constructor but there's no need for clients to access it, while clients + // should be able to access this class and its getters. + @Deprecated(since = "8.3.0", forRemoval = true) + public IncrementalWriteConfig(String hashKeyName, String timestampKeyName, + boolean canonicalizeJson, Consumer skippedDocumentsConsumer, String[] jsonExclusions, String[] xmlExclusions, Map xmlNamespaces, String schemaName, String viewName) { + this(hashKeyName, null, timestampKeyName, canonicalizeJson, skippedDocumentsConsumer, jsonExclusions, xmlExclusions, xmlNamespaces, schemaName, viewName); + } + + IncrementalWriteConfig(String hashKeyName, String sourceUriKeyName, String timestampKeyName, + boolean canonicalizeJson, + Consumer skippedDocumentsConsumer, + String[] jsonExclusions, String[] xmlExclusions, Map xmlNamespaces, + String schemaName, String viewName) { this.hashKeyName = hashKeyName; + this.sourceUriKeyName = sourceUriKeyName; this.timestampKeyName = timestampKeyName; this.canonicalizeJson = canonicalizeJson; this.skippedDocumentsConsumer = skippedDocumentsConsumer; @@ -45,6 +59,13 @@ public String getHashKeyName() { return hashKeyName; } + /** + * @since 8.3.0 + */ + public String getSourceUriKeyName() { + return sourceUriKeyName; + } + public String getTimestampKeyName() { return timestampKeyName; } diff --git a/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteFilter.java b/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteFilter.java index 25be19f06..31625eba4 100644 --- a/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteFilter.java +++ b/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteFilter.java @@ -46,6 +46,7 @@ public static Builder newBuilder() { public static class Builder { private String hashKeyName = "incrementalWriteHash"; + private String sourceUriKeyName; private String timestampKeyName; private boolean canonicalizeJson = true; private Consumer skippedDocumentsConsumer; @@ -75,6 +76,26 @@ public Builder timestampKeyName(String keyName) { return this; } + /** + * Configures this filter to match existing documents by a stable "source URI" metadata + * value instead of the document's actual/current URI. Useful when a document may be + * relocated to a different URI after this filter writes it, since matching by URI alone + * would otherwise stop working. If set, the filter will expect a field range index to exist on the + * specified metadata key. The field range index will be utilized to find matching documents instead of a + * document query on URIs. + * + * @param keyName the name of the MarkLogic metadata key that will hold each document's + * source URI; defaults to null, which means matching is based on the + * document's actual URI, as it has always worked. + * @since 8.3.0 + */ + public Builder sourceUriKeyName(String keyName) { + if (keyName != null && !keyName.trim().isEmpty()) { + this.sourceUriKeyName = keyName; + } + return this; + } + /** * @param canonicalizeJson whether to canonicalize JSON content before hashing; defaults to true. * Delegates to https://github.com/erdtman/java-json-canonicalization for canonicalization. @@ -146,8 +167,9 @@ public Builder fromView(String schemaName, String viewName) { public IncrementalWriteFilter build() { validateJsonExclusions(); validateXmlExclusions(); - IncrementalWriteConfig config = new IncrementalWriteConfig(hashKeyName, timestampKeyName, canonicalizeJson, - skippedDocumentsConsumer, jsonExclusions, xmlExclusions, xmlNamespaces, schemaName, viewName); + IncrementalWriteConfig config = new IncrementalWriteConfig(hashKeyName, sourceUriKeyName, timestampKeyName, + canonicalizeJson, skippedDocumentsConsumer, jsonExclusions, xmlExclusions, xmlNamespaces, + schemaName, viewName); if (schemaName != null && viewName != null) { return new IncrementalWriteFromViewFilter(config); @@ -240,14 +262,16 @@ protected final DocumentWriteSet filterDocuments(Context context, Function existingHashes = new RowTemplate(context.getDatabaseClient()).query(op -> op.fromLexicons(Map.of( - "uri", op.cts.uriReference(), + "uri", sourceUriKeyName != null + ? op.cts.fieldReference(sourceUriKeyName) + : op.cts.uriReference(), "hash", op.cts.fieldReference(getConfig().getHashKeyName()) )).where( - op.cts.documentQuery(op.xs.stringSeq(uris)) + sourceUriKeyName != null + ? op.cts.fieldRangeQuery(op.xs.string(sourceUriKeyName), op.xs.string("="), op.xs.stringSeq(uris)) + : op.cts.documentQuery(op.xs.stringSeq(uris)) ), rows -> { diff --git a/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteFromViewFilter.java b/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteFromViewFilter.java index 754944b5b..5028ece50 100644 --- a/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteFromViewFilter.java +++ b/marklogic-client-api/src/main/java/com/marklogic/client/datamovement/filter/IncrementalWriteFromViewFilter.java @@ -5,6 +5,7 @@ import com.marklogic.client.FailedRequestException; import com.marklogic.client.document.DocumentWriteSet; +import com.marklogic.client.expression.PlanBuilder; import com.marklogic.client.row.RowTemplate; import java.util.HashMap; @@ -28,10 +29,16 @@ public DocumentWriteSet apply(Context context) { final String[] uris = getUrisInBatch(context.getDocumentWriteSet()); try { - Map existingHashes = new RowTemplate(context.getDatabaseClient()).query(op -> - op.fromView(getConfig().getSchemaName(), getConfig().getViewName(), "") - .where(op.cts.documentQuery(op.xs.stringSeq(uris))) - , + Map existingHashes = new RowTemplate(context.getDatabaseClient()).query(op -> { + PlanBuilder.ModifyPlan plan = op.fromView(getConfig().getSchemaName(), getConfig().getViewName(), ""); + final String sourceUriName = getConfig().getSourceUriKeyName(); + if (sourceUriName != null && !sourceUriName.trim().isEmpty()) { + plan = plan.where(op.cts.fieldRangeQuery(op.xs.string(sourceUriName), op.xs.string("="), op.xs.stringSeq(uris))); + } else { + plan = plan.where(op.cts.documentQuery(op.xs.stringSeq(uris))); + } + return plan; + }, rows -> { Map map = new HashMap<>(); rows.forEach(row -> { diff --git a/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteFilterTest.java b/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteFilterTest.java index bdaab89b4..d982d2723 100644 --- a/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteFilterTest.java +++ b/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteFilterTest.java @@ -49,4 +49,38 @@ void addHashToMetadata() { assertEquals("12345", metadata2.getMetadataValues().get("theField"), "hash field should be added"); assertEquals(timestamp, metadata2.getMetadataValues().get("theTimestamp"), "timestamp should be added"); } + + @Test + void addHashToMetadataStampsSourceUriWhenConfigured() { + DocumentWriteOperation doc = new DocumentWriteOperationImpl("/source/1.json", null, new StringHandle("{}")); + + doc = IncrementalWriteFilter.addHashToMetadata(doc, "theField", "sourceUri", 12345, null, null); + + DocumentMetadataHandle metadata = (DocumentMetadataHandle) doc.getMetadata(); + assertEquals("/source/1.json", metadata.getMetadataValues().get("sourceUri"), + "the document's own URI should be stamped as the source URI when none was already set"); + } + + @Test + void addHashToMetadataPreservesExplicitSourceUri() { + DocumentMetadataHandle metadata = new DocumentMetadataHandle() + .withMetadataValue("sourceUri", "/original/source.json"); + DocumentWriteOperation doc = new DocumentWriteOperationImpl("/relocated/1.json", metadata, new StringHandle("{}")); + + doc = IncrementalWriteFilter.addHashToMetadata(doc, "theField", "sourceUri", 12345, null, null); + + DocumentMetadataHandle newMetadata = (DocumentMetadataHandle) doc.getMetadata(); + assertEquals("/original/source.json", newMetadata.getMetadataValues().get("sourceUri"), + "an explicit source URI already present in metadata should not be overwritten"); + } + + @Test + void addHashToMetadataWithoutSourceUriKeyNameDoesNotStampAnything() { + DocumentWriteOperation doc = new DocumentWriteOperationImpl("/1.json", null, new StringHandle("{}")); + + doc = IncrementalWriteFilter.addHashToMetadata(doc, "theField", null, 12345, null, null); + + DocumentMetadataHandle metadata = (DocumentMetadataHandle) doc.getMetadata(); + assertEquals(1, metadata.getMetadataValues().size(), "only the hash field should be present"); + } } diff --git a/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteSourceUriKeyNameTest.java b/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteSourceUriKeyNameTest.java new file mode 100644 index 000000000..3e0795098 --- /dev/null +++ b/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteSourceUriKeyNameTest.java @@ -0,0 +1,107 @@ +/* + * Copyright (c) 2010-2026 Progress Software Corporation and/or its subsidiaries or affiliates. All Rights Reserved. + */ +package com.marklogic.client.datamovement.filter; + +import com.marklogic.client.document.DocumentWriteOperation; +import com.marklogic.client.impl.DocumentWriteOperationImpl; +import com.marklogic.client.io.DocumentMetadataHandle; +import com.marklogic.client.io.StringHandle; +import com.marklogic.client.test.Common; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class IncrementalWriteSourceUriKeyNameTest extends AbstractIncrementalWriteTest { + + private static final String SOURCE_URI_KEY_NAME = "incrementalWriteSourceUri"; + private static final String SOURCE_URI = "/incremental/test/source-uri-doc.xml"; + private static final String RELOCATED_URI = "/incremental/test/relocated-doc.xml"; + + @Override + @BeforeEach + void setup() { + super.setup(); + filter = IncrementalWriteFilter.newBuilder() + .sourceUriKeyName(SOURCE_URI_KEY_NAME) + .onDocumentsSkipped(docs -> skippedCount.addAndGet(docs.length)) + .build(); + } + + @Test + void stampsSourceUriMetadataAutomatically() { + writeDoc(SOURCE_URI, "original content"); + + DocumentMetadataHandle metadata = Common.client.newDocumentManager().readMetadata(SOURCE_URI, + new DocumentMetadataHandle()); + assertEquals(SOURCE_URI, metadata.getMetadataValues().get(SOURCE_URI_KEY_NAME), + "The filter should automatically stamp the configured metadata key with the document's own URI."); + } + + @Test + void matchesAfterExternalRelocation() { + // Initial write via the filter. + writeDoc(SOURCE_URI, "original content"); + assertEquals(1, writtenCount.get()); + assertEquals(0, skippedCount.get()); + + relocateDocument(); + + // Same source URI, unchanged content - should be recognized as unchanged even though the document + // now physically lives at a different URI, because the lookup matches on the source-URI field value. + writeDoc(SOURCE_URI, "original content"); + assertEquals(1, writtenCount.get(), "No new write should have occurred since the document was skipped."); + assertEquals(1, skippedCount.get()); + + assertNull(Common.client.newDocumentManager().exists(SOURCE_URI), + "Nothing should have been written back to the source URI since the write was skipped."); + assertNotNull(Common.client.newDocumentManager().exists(RELOCATED_URI), + "The relocated document should still be the only copy of this content."); + } + + @Test + void changedContentIsWrittenAfterExternalRelocation() { + writeDoc(SOURCE_URI, "original content"); + relocateDocument(); + + // Content has changed, so the write should go through - to the source URI, not the relocated one. + writeDoc(SOURCE_URI, "modified content"); + assertEquals(2, writtenCount.get()); + assertEquals(0, skippedCount.get()); + + String newContent = Common.client.newTextDocumentManager().read(SOURCE_URI, new StringHandle()).get(); + assertNotNull(newContent); + assertTrue(newContent.contains("modified content")); + + DocumentMetadataHandle metadata = Common.client.newDocumentManager().readMetadata(SOURCE_URI, + new DocumentMetadataHandle()); + assertEquals(SOURCE_URI, metadata.getMetadataValues().get(SOURCE_URI_KEY_NAME), + "The rewritten document should have the source-URI metadata value stamped again."); + } + + private void writeDoc(String uri, String content) { + List ops = new ArrayList<>(); + ops.add(new DocumentWriteOperationImpl(uri, METADATA, new StringHandle(content))); + writeDocs(ops); + } + + /** + * Simulates a process outside of this filter relocating the document to a new URI while preserving its + * metadata (in particular, the source-URI metadata value), then removing the original. + */ + private void relocateDocument() { + DocumentMetadataHandle metadata = Common.client.newDocumentManager().readMetadata(SOURCE_URI, + new DocumentMetadataHandle()); + String content = Common.client.newTextDocumentManager().read(SOURCE_URI, new StringHandle()).get(); + + Common.client.newTextDocumentManager().write(RELOCATED_URI, metadata, new StringHandle(content)); + Common.client.newDocumentManager().delete(SOURCE_URI); + } +} diff --git a/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteTest.java b/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteTest.java index 4c89e3dc5..c80f3e405 100644 --- a/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteTest.java +++ b/marklogic-client-api/src/test/java/com/marklogic/client/datamovement/filter/IncrementalWriteTest.java @@ -227,6 +227,38 @@ void fromView() { verifyIncrementalWriteWorks(); } + @Test + void fromViewWithCustomSourceUri() { + filter = IncrementalWriteFilter.newBuilder() + .fromView("javaClient", "incrementalWriteHash") + .timestampKeyName("incrementalWriteTimestamp") + .sourceUriKeyName("incrementalWriteSourceUri") + .onDocumentsSkipped(docs -> skippedCount.addAndGet(docs.length)) + .build(); + + verifyIncrementalWriteWorks(); + verifyDocumentsHaveSourceUriInMetadataKey(); + } + + @Test + void fromViewWithInvalidCustomSourceUri() { + filter = IncrementalWriteFilter.newBuilder() + .fromView("javaClient", "incrementalWriteHash") + .timestampKeyName("incrementalWriteTimestamp") + .sourceUriKeyName("noFieldRangeIndexOnThis") + .onDocumentsSkipped(docs -> skippedCount.addAndGet(docs.length)) + .build(); + + writeTenDocuments(); + assertNotNull(batchFailure.get()); + + String message = batchFailure.get().getMessage(); + assertTrue(message.contains("Field not defined: noFieldRangeIndexOnThis"), + "This test configures a sourceUriKeyName using a metadata key that does not have a field range index " + + "on it. This should cause an error as the Optic query against the view is expected to constrain on " + + "documents via a field range query. Actual message: " + message); + } + @Test void loadWithViewFilterThenVerifyLexiconFilterSkipsAll() { filter = IncrementalWriteFilter.newBuilder() @@ -364,6 +396,19 @@ private void verifyDocumentsHasHashInMetadataKey() { } } + private void verifyDocumentsHaveSourceUriInMetadataKey() { + GenericDocumentManager mgr = Common.client.newDocumentManager(); + mgr.setMetadataCategories(DocumentManager.Metadata.METADATAVALUES); + DocumentPage page = mgr.search(Common.client.newQueryManager().newStructuredQueryBuilder().collection("incremental-test"), 1); + while (page.hasNext()) { + DocumentRecord doc = page.next(); + DocumentMetadataHandle metadata = doc.getMetadata(new DocumentMetadataHandle()); + + String sourceUri = metadata.getMetadataValues().get("incrementalWriteSourceUri"); + assertEquals(doc.getUri(), sourceUri, "Document " + doc.getUri() + " should have an incrementalWriteSourceUri value equal to its own URI."); + } + } + private void modifyFiveDocuments() { docs = new ArrayList<>(); for (int i = 6; i <= 10; i++) { diff --git a/test-app/src/main/ml-config/databases/content-database.json b/test-app/src/main/ml-config/databases/content-database.json index 965bb45d3..87afbb02f 100644 --- a/test-app/src/main/ml-config/databases/content-database.json +++ b/test-app/src/main/ml-config/databases/content-database.json @@ -198,6 +198,15 @@ "fast-case-sensitive-searches": false, "fast-diacritic-sensitive-searches": false }, + { + "field-name": "incrementalWriteSourceUri", + "metadata": "", + "stemmed-searches": "off", + "word-searches": false, + "fast-phrase-searches": false, + "fast-case-sensitive-searches": false, + "fast-diacritic-sensitive-searches": false + }, { "field-name": "myWriteHash", "metadata": "", @@ -229,6 +238,13 @@ "range-value-positions": false, "invalid-values": "reject" }, + { + "scalar-type": "string", + "field-name": "incrementalWriteSourceUri", + "collation": "http://marklogic.com/collation/", + "range-value-positions": false, + "invalid-values": "reject" + }, { "scalar-type": "unsignedLong", "field-name": "myWriteHash",