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
14 changes: 14 additions & 0 deletions marklogic-client-api/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
*/

plugins {
id 'net.saliman.properties' version '1.6.0'
id 'maven-publish'
}

Expand Down Expand Up @@ -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")) {
Comment thread
rjrudin marked this conversation as resolved.
maven {
name = "GitHubPackages"
url = uri(ghPackagesUrl)
credentials {
username = ghActor
password = ghToken
}
}
}

maven {
if (project.hasProperty("mavenUser")) {
credentials {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<DocumentWriteOperation[]> skippedDocumentsConsumer;
Expand All @@ -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<DocumentWriteOperation[]> skippedDocumentsConsumer,
String[] jsonExclusions, String[] xmlExclusions, Map<String, String> 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<DocumentWriteOperation[]> skippedDocumentsConsumer,
String[] jsonExclusions, String[] xmlExclusions, Map<String, String> xmlNamespaces,
String schemaName, String viewName) {
this.hashKeyName = hashKeyName;
this.sourceUriKeyName = sourceUriKeyName;
this.timestampKeyName = timestampKeyName;
this.canonicalizeJson = canonicalizeJson;
this.skippedDocumentsConsumer = skippedDocumentsConsumer;
Expand All @@ -45,6 +59,13 @@ public String getHashKeyName() {
return hashKeyName;
}

/**
* @since 8.3.0
*/
public String getSourceUriKeyName() {
return sourceUriKeyName;
}

public String getTimestampKeyName() {
return timestampKeyName;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<DocumentWriteOperation[]> skippedDocumentsConsumer;
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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);
Comment thread
rjrudin marked this conversation as resolved.

if (schemaName != null && viewName != null) {
return new IncrementalWriteFromViewFilter(config);
Expand Down Expand Up @@ -240,14 +262,16 @@ protected final DocumentWriteSet filterDocuments(Context context, Function<Strin

if (existingHash != null) {
if (!existingHash.equals(contentHash)) {
newWriteSet.add(addHashToMetadata(doc, config.getHashKeyName(), contentHash, config.getTimestampKeyName(), timestamp));
newWriteSet.add(addHashToMetadata(doc, config.getHashKeyName(), config.getSourceUriKeyName(),
contentHash, config.getTimestampKeyName(), timestamp));
} else if (config.getSkippedDocumentsConsumer() != null) {
skippedDocuments.add(doc);
} else {
// No consumer, so skip the document silently.
}
} else {
newWriteSet.add(addHashToMetadata(doc, config.getHashKeyName(), contentHash, config.getTimestampKeyName(), timestamp));
newWriteSet.add(addHashToMetadata(doc, config.getHashKeyName(), config.getSourceUriKeyName(),
contentHash, config.getTimestampKeyName(), timestamp));
}
}

Expand Down Expand Up @@ -310,6 +334,24 @@ private long computeHash(String content) {

protected static DocumentWriteOperation addHashToMetadata(DocumentWriteOperation op, String hashKeyName, long hash,
String timestampKeyName, String timestamp) {
return addHashToMetadata(op, hashKeyName, null, hash, timestampKeyName, timestamp);
}

/**
* @param op the write operation to add hash (and optionally source-URI) metadata to.
* @param hashKeyName the metadata key that will hold the content hash.
* @param sourceUriKeyName the metadata key that will hold the document's source URI, or null if this
* feature is not in use. If the operation's existing metadata already has a
* value for this key, it is preserved rather than overwritten with
* {@code op.getUri()}.
* @param hash the newly computed content hash.
* @param timestampKeyName the metadata key that will hold the write timestamp, or null if timestamps
* are not in use.
* @param timestamp the timestamp to store, if {@code timestampKeyName} is not null.
*/
protected static DocumentWriteOperation addHashToMetadata(DocumentWriteOperation op, String hashKeyName,
String sourceUriKeyName, long hash,
String timestampKeyName, String timestamp) {
DocumentMetadataHandle newMetadata = new DocumentMetadataHandle();
if (op.getMetadata() != null) {
DocumentMetadataHandle originalMetadata = (DocumentMetadataHandle) op.getMetadata();
Expand All @@ -324,6 +366,12 @@ protected static DocumentWriteOperation addHashToMetadata(DocumentWriteOperation
if (timestampKeyName != null && !timestampKeyName.trim().isEmpty()) {
newMetadata.getMetadataValues().put(timestampKeyName, timestamp);
}
if (sourceUriKeyName != null && !sourceUriKeyName.trim().isEmpty()
&& !newMetadata.getMetadataValues().containsKey(sourceUriKeyName)) {
// Only stamp this automatically if the caller hasn't already supplied an explicit value -
// e.g. via metadata set directly on the incoming write operation.
newMetadata.getMetadataValues().put(sourceUriKeyName, op.getUri());
Comment on lines +369 to +373
}

return new DocumentWriteOperationImpl(op.getUri(), newMetadata, op.getContent(), op.getTemporalDocumentURI());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,10 @@

/**
* Uses an Optic fromLexicons query that depends on a field range index to retrieve URIs and
* hash values.
* hash values. When {@code sourceUriKeyName} is configured (see
* {@link IncrementalWriteFilter.Builder#sourceUriKeyName(String)}), the "URI" side of the
* lookup is instead backed by a field over that metadata key, since the matched document may
* have since been relocated to a different physical URI.
*
* @since 8.1.0
*/
Expand All @@ -26,14 +29,19 @@ class IncrementalWriteFromLexiconsFilter extends IncrementalWriteFilter {
@Override
public DocumentWriteSet apply(Context context) {
final String[] uris = getUrisInBatch(context.getDocumentWriteSet());
final String sourceUriKeyName = getConfig().getSourceUriKeyName();

try {
Map<String, Long> 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 -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -28,10 +29,16 @@ public DocumentWriteSet apply(Context context) {
final String[] uris = getUrisInBatch(context.getDocumentWriteSet());

try {
Map<String, Long> existingHashes = new RowTemplate(context.getDatabaseClient()).query(op ->
op.fromView(getConfig().getSchemaName(), getConfig().getViewName(), "")
.where(op.cts.documentQuery(op.xs.stringSeq(uris)))
,
Map<String, Long> 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)));
Comment on lines +35 to +36
} else {
plan = plan.where(op.cts.documentQuery(op.xs.stringSeq(uris)));
}
return plan;
},
rows -> {
Map<String, Long> map = new HashMap<>();
rows.forEach(row -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
}
Original file line number Diff line number Diff line change
@@ -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, "<doc>original content</doc>");

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, "<doc>original content</doc>");
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, "<doc>original content</doc>");
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, "<doc>original content</doc>");
relocateDocument();

// Content has changed, so the write should go through - to the source URI, not the relocated one.
writeDoc(SOURCE_URI, "<doc>modified content</doc>");
assertEquals(2, writtenCount.get());
assertEquals(0, skippedCount.get());

String newContent = Common.client.newTextDocumentManager().read(SOURCE_URI, new StringHandle()).get();
assertNotNull(newContent);
assertTrue(newContent.contains("<doc>modified content</doc>"));

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<DocumentWriteOperation> 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);
}
}
Loading
Loading