From a4f7ae1c4e75b31b170247a1ae2476fe7773a74f Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Thu, 24 Sep 2026 18:27:28 +0530 Subject: [PATCH 1/9] fix(server): converge schema caches across HStore servers A schema change made through one Server could stay invisible on the others indefinitely: updates and status flips publish no PD event, and the PD watch does not replay events lost while it reconnects. Every schema change now writes a per-graph opaque version to PD at HUGEGRAPH/{cluster}/SCHEMA_VERSION/{graphspace}/{graph}. Each Server checks it once per schema.sync.reconcile_interval (10s by default) and clears that graph's schema cache when it changed, so a change becomes visible everywhere within about one interval without relying on events. A burst of changes costs each Server at most one clear per interval. Existing events are unchanged. A clear caused by a remote change and a reader's cache update are now exclusive on the graph's SchemaCaches, and a reader skips its update if such a clear ran after it read, so data read before the clear can't be cached after it. New options: schema.sync.enabled (default true) and schema.sync.reconcile_interval (0 to 3600 seconds, 0 disables checks). Fixes #3235 --- .../apache/hugegraph/core/GraphManager.java | 7 + .../apache/hugegraph/StandardHugeGraph.java | 15 + .../cache/CachedSchemaTransactionV2.java | 218 +++++++++-- .../cache/SchemaVersionReconciler.java | 231 ++++++++++++ .../apache/hugegraph/config/CoreOptions.java | 25 ++ .../apache/hugegraph/meta/MetaManager.java | 14 + .../meta/managers/GraphMetaManager.java | 32 ++ .../static/conf/graphs/hugegraph.properties | 4 + .../apache/hugegraph/unit/UnitTestSuite.java | 4 + .../cache/CachedSchemaTransactionTest.java | 284 +++++++++++++- .../cache/SchemaVersionReconcilerTest.java | 355 ++++++++++++++++++ .../unit/core/GraphManagerDropGraphTest.java | 124 ++++++ 12 files changed, 1271 insertions(+), 42 deletions(-) create mode 100644 hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/SchemaVersionReconciler.java create mode 100644 hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/SchemaVersionReconcilerTest.java create mode 100644 hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/GraphManagerDropGraphTest.java diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java index a386a0ca67..d0122ca7a2 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java @@ -2349,6 +2349,13 @@ public void dropGraph(String graphSpace, String name, boolean clear) { } g.clearBackend(); + try { + // The schema version kept in PD for the schema cache sync + this.metaManager.deleteSchemaVersion(graphSpace, name); + } catch (Exception e) { + LOG.warn("Failed to delete the schema version of graph {}", + graphName, e); + } try { g.close(); } catch (Exception e) { diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java index 8f2ff95d01..da39ae1046 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java @@ -1131,6 +1131,10 @@ public synchronized void close() throws Exception { this.closeTx(); } finally { this.closed = true; + if (this.isHstore()) { + // After closed is set, so no new reconciler can be created + CachedSchemaTransactionV2.stopReconciler(this.spaceGraphName()); + } this.storeProvider.close(); LockUtil.destroy(this.spaceGraphName()); } @@ -1161,6 +1165,17 @@ public void create(String configPath, GlobalMasterInfo nodeInfo) { @Override public void drop() { this.clearBackend(); + // Not this.option(): it only allows the options in ALLOWED_CONFIGS + if (this.isHstore() && + this.configuration().get(CoreOptions.SCHEMA_SYNC_ENABLED)) { + try { + MetaManager.instance().deleteSchemaVersion(this.graphSpace(), + this.name()); + } catch (Exception e) { + LOG.warn("Failed to delete the schema version of graph {}: {}", + this.spaceGraphName(), e.toString()); + } + } HugeConfig config = this.configuration(); this.storeProvider.onDeleteConfig(config); diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java index 74986b2e99..485368c952 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java @@ -23,15 +23,19 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import java.util.function.Consumer; +import org.apache.hugegraph.HugeGraph; import org.apache.hugegraph.HugeGraphParams; import org.apache.hugegraph.backend.id.Id; import org.apache.hugegraph.backend.id.IdGenerator; import org.apache.hugegraph.backend.store.ram.IntObjectMap; import org.apache.hugegraph.backend.tx.SchemaTransactionV2; import org.apache.hugegraph.config.CoreOptions; +import org.apache.hugegraph.config.HugeConfig; import org.apache.hugegraph.event.EventHub; import org.apache.hugegraph.event.EventListener; import org.apache.hugegraph.meta.MetaDriver; @@ -76,11 +80,23 @@ public class CachedSchemaTransactionV2 extends SchemaTransactionV2 { private static final String SCHEMA_CACHE_CLEAR_SOURCE = UUID.randomUUID().toString(); + /* + * One reconciler per open graph, keyed by space graph name. Created by the + * first schema transaction of the graph and stopped by the graph close, + * not by a transaction close: schema transactions are per thread and a + * thread's transaction is closed after every task. + */ + private static final ConcurrentMap + RECONCILERS = new ConcurrentHashMap<>(); + private final Cache idCache; private final Cache nameCache; private final SchemaCaches arrayCaches; + // Null if schema.sync.enabled is false + private final SchemaVersionReconciler reconciler; + private EventListener storeEventListener; private EventListener cacheEventListener; @@ -101,6 +117,8 @@ public CachedSchemaTransactionV2(MetaDriver metaDriver, } this.arrayCaches = attachment; this.listenChanges(); + + this.reconciler = ensureReconciler(graphParams); } private static Id generateId(HugeType type, Id id) { @@ -131,19 +149,39 @@ private static String cacheName(String prefix, String spaceGraphName) { private static void clearSchemaCache(String spaceGraphName) { Map> caches = CacheManager.instance().caches(); + Cache nameCache = caches.get(cacheName(NAME_CACHE_PREFIX, + spaceGraphName)); + Cache idCache = caches.get(cacheName(ID_CACHE_PREFIX, + spaceGraphName)); + SchemaCaches arrayCaches = idCache == null ? + null : idCache.attachment(); + if (arrayCaches == null) { + clearSchemaCache(nameCache, null, idCache); + return; + } + /* + * This clear is caused by a schema change on another server, so data + * a reader got before it may be stale: the new generation, set after + * the wipe, makes such a reader skip its cache update. The lock makes + * the clear and a reader's cache update exclusive. + */ + synchronized (arrayCaches) { + clearSchemaCache(nameCache, arrayCaches, idCache); + arrayCaches.nextGeneration(); + } + } + + private static void clearSchemaCache(Cache nameCache, + SchemaCaches arrayCaches, + Cache idCache) { // Clear name cache first so the (name -> id -> object) lookup path // fails fast instead of returning a stale object backed by an // already-empty id cache during the TOCTOU window. - Cache nameCache = caches.get(cacheName(NAME_CACHE_PREFIX, - spaceGraphName)); if (nameCache != null) { nameCache.clear(); } - Cache idCache = caches.get(cacheName(ID_CACHE_PREFIX, - spaceGraphName)); if (idCache != null) { - SchemaCaches arrayCaches = idCache.attachment(); if (arrayCaches != null) { arrayCaches.clear(); } @@ -151,6 +189,46 @@ private static void clearSchemaCache(String spaceGraphName) { } } + private static SchemaVersionReconciler ensureReconciler( + HugeGraphParams params) { + HugeConfig config = params.configuration(); + if (!config.get(CoreOptions.SCHEMA_SYNC_ENABLED)) { + return null; + } + long intervalMs = 1000L * config.get( + CoreOptions.SCHEMA_SYNC_RECONCILE_INTERVAL); + HugeGraph graph = params.graph(); + SchemaVersionReconciler reconciler = RECONCILERS.compute( + graph.spaceGraphName(), (name, existing) -> { + if (existing != null && !existing.stopped()) { + return existing; + } + if (params.closed()) { + // Don't leave a reconciler of a closed graph in the map + return null; + } + SchemaVersionReconciler created = new SchemaVersionReconciler( + graph.graphSpace(), graph.name(), + SCHEMA_CACHE_CLEAR_SOURCE, new MetaVersionStore(), + () -> clearSchemaCache(name), params::closed); + if (intervalMs > 0L) { + created.start(intervalMs); + } + return created; + }); + if (reconciler != null && intervalMs > 0L) { + reconciler.ensureScheduled(intervalMs); + } + return reconciler; + } + + public static void stopReconciler(String spaceGraphName) { + SchemaVersionReconciler reconciler = RECONCILERS.remove(spaceGraphName); + if (reconciler != null) { + reconciler.stop(); + } + } + private void listenChanges() { // Listen store event: "store.init", "store.clear", ... Set storeEvents = ImmutableSet.of(Events.STORE_INIT, @@ -258,13 +336,17 @@ static void handleSchemaCacheClearEvent(T response) { } public void clearCache(boolean notify) { - // Same TOCTOU ordering as clearSchemaCache(String): clear nameCache - // first, then the array attachment, then idCache last. - this.nameCache.clear(); - this.arrayCaches.clear(); - this.idCache.clear(); + // Exclusive with a reader's cache update, see getAllSchema() + synchronized (this.arrayCaches) { + // Same TOCTOU ordering as clearSchemaCache(String): clear nameCache + // first, then the array attachment, then idCache last. + this.nameCache.clear(); + this.arrayCaches.clear(); + this.idCache.clear(); + } if (notify) { + this.bumpSchemaVersion(); this.maybeNotifySchemaCacheClear(); } } @@ -318,8 +400,10 @@ protected void updateSchema(SchemaElement schema, super.updateSchema(schema, updateCallback); this.updateCache(schema); - // Status transitions are internal bookkeeping; notifying here causes a - // broadcast storm for every updateSchemaStatus() call from background jobs. + // No meta event here: one per updateSchemaStatus() call from background + // jobs would be a broadcast storm. The schema version has no fan-out, + // other servers check it at most once per reconcile interval. + this.bumpSchemaVersion(); } @Override @@ -328,6 +412,7 @@ protected void addSchema(SchemaElement schema) { this.updateCache(schema); + this.bumpSchemaVersion(); // Schema additions must always propagate to remote nodes regardless // of TASK_SYNC_DELETION (which only gates removal flows). this.notifySchemaCacheClear(); @@ -354,9 +439,16 @@ public void removeSchema(SchemaElement schema) { this.invalidateCache(schema.type(), schema.id()); + this.bumpSchemaVersion(); this.maybeNotifySchemaCacheClear(); } + private void bumpSchemaVersion() { + if (this.reconciler != null) { + this.reconciler.bump(); + } + } + private void maybeNotifySchemaCacheClear() { // Only suppress notifications for removal tasks when // TASK_SYNC_DELETION=true: the caller propagates cache invalidation @@ -384,25 +476,40 @@ protected T getSchema(HugeType type, Id id) { } } + long generation = this.arrayCaches.generation(); Id prefixedId = generateId(type, id); Object value = this.idCache.get(prefixedId); - if (value == null) { - value = super.getSchema(type, id); - if (value != null) { - this.resetCachedAllIfReachedCapacity(); - - this.idCache.update(prefixedId, value); - - SchemaElement schema = (SchemaElement) value; - Id prefixedName = generateId(schema.type(), schema.name()); - this.nameCache.update(prefixedName, schema); + if (value != null) { + if (this.arrayCaches.fits(id)) { + synchronized (this.arrayCaches) { + // Don't promote what a remote change cleared meanwhile + if (generation == this.arrayCaches.generation()) { + // update optimized array cache + this.arrayCaches.updateIfNeeded((SchemaElement) value); + } + } } + return (T) value; } - // update optimized array cache - this.arrayCaches.updateIfNeeded((SchemaElement) value); + SchemaElement schema = super.getSchema(type, id); + if (schema != null) { + synchronized (this.arrayCaches) { + // Don't cache what was read before a remote change cleared + if (generation == this.arrayCaches.generation()) { + this.resetCachedAllIfReachedCapacity(); - return (T) value; + this.idCache.update(prefixedId, schema); + + Id prefixedName = generateId(schema.type(), schema.name()); + this.nameCache.update(prefixedName, schema); + + // update optimized array cache + this.arrayCaches.updateIfNeeded(schema); + } + } + } + return (T) schema; } @Override @@ -438,18 +545,31 @@ protected List getAllSchema(HugeType type) { }); return results; } else { + long generation = this.arrayCaches.generation(); results = super.getAllSchema(type); long free = this.idCache.capacity() - this.idCache.size(); if (results.size() <= free) { - // Update cache - for (T schema : results) { - Id prefixedId = generateId(schema.type(), schema.id()); - this.idCache.update(prefixedId, schema); - - Id prefixedName = generateId(schema.type(), schema.name()); - this.nameCache.update(prefixedName, schema); + /* + * Under the lock a clear can't wipe the caches halfway through + * this update, which would leave cachedAll set over a partial + * id cache. Skip the update if a remote change cleared the + * caches after the storage read. + */ + synchronized (this.arrayCaches) { + if (generation == this.arrayCaches.generation()) { + // Update cache + for (T schema : results) { + Id prefixedId = generateId(schema.type(), + schema.id()); + this.idCache.update(prefixedId, schema); + + Id prefixedName = generateId(schema.type(), + schema.name()); + this.nameCache.update(prefixedName, schema); + } + this.cachedTypes().putIfAbsent(type, true); + } } - this.cachedTypes().putIfAbsent(type, true); } return results; } @@ -467,6 +587,8 @@ public void clear() { // Clear schema info firstly super.clear(); this.clearCache(false); + // Write a new version instead of deleting it: "" isn't unique + this.bumpSchemaVersion(); this.notifySchemaCacheClear(); } @@ -481,6 +603,9 @@ private static final class SchemaCaches { private final CachedTypes cachedTypes; + // Incremented by every clear caused by a remote schema change + private final AtomicLong generation; + public SchemaCaches(int size) { // TODO: improve size of each type for optimized array cache this.size = size; @@ -491,6 +616,19 @@ public SchemaCaches(int size) { this.ils = new IntObjectMap<>(size); this.cachedTypes = new CachedTypes(); + this.generation = new AtomicLong(0L); + } + + public long generation() { + return this.generation.get(); + } + + public void nextGeneration() { + this.generation.incrementAndGet(); + } + + public boolean fits(Id id) { + return id.number() && id.asLong() > 0L && id.asLong() < this.size; } public void updateIfNeeded(V schema) { @@ -608,4 +746,18 @@ private static class CachedTypes private static final long serialVersionUID = -2215549791679355996L; } + + private static final class MetaVersionStore + implements SchemaVersionReconciler.VersionStore { + + @Override + public String read(String graphSpace, String graph) { + return MetaManager.instance().getSchemaVersion(graphSpace, graph); + } + + @Override + public void write(String graphSpace, String graph, String version) { + MetaManager.instance().putSchemaVersion(graphSpace, graph, version); + } + } } diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/SchemaVersionReconciler.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/SchemaVersionReconciler.java new file mode 100644 index 0000000000..f6b3d2bd14 --- /dev/null +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/SchemaVersionReconciler.java @@ -0,0 +1,231 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with this + * work for additional information regarding copyright ownership. The ASF + * licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hugegraph.backend.cache; + +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.BooleanSupplier; + +import org.apache.hugegraph.util.E; +import org.apache.hugegraph.util.Log; +import org.slf4j.Logger; + +/** + * Keeps the schema cache of one graph on this server consistent with schema + * changes made through other servers. + *

+ * Every schema change writes a new opaque version of the graph to PD + * ({@link #bump()}). Each tick ({@link #run()}) reads the version and, if it + * differs from the one adopted last time, clears the schema cache of the + * graph and adopts the version it read before the clear. Versions are only + * compared for equality. The PD watch events are not involved, so a lost + * event or a dead watch can delay a change by at most one interval. + *

+ * Changes made through this server are not skipped: a writer could otherwise + * adopt its own version after another server wrote a newer schema element, + * and keep that element stale. + */ +public final class SchemaVersionReconciler implements Runnable { + + private static final Logger LOG = Log.logger(SchemaVersionReconciler.class); + + private static final AtomicLong VERSION_SEQ = new AtomicLong(); + private static final String NONE = ""; + + // One daemon thread runs the ticks of every graph in this JVM + private static final ScheduledExecutorService SCHEDULER = newScheduler(); + + private final String graphSpace; + private final String graph; + private final String spaceGraphName; + private final String source; + private final VersionStore store; + private final Runnable clearCache; + private final BooleanSupplier graphClosed; + + private final AtomicBoolean pendingWrite; + // The version adopted by the last clear, null before the first tick + private volatile String applied; + private volatile boolean unreachable; + private volatile boolean stopped; + private ScheduledFuture future; + + public SchemaVersionReconciler(String graphSpace, String graph, + String source, VersionStore store, + Runnable clearCache, + BooleanSupplier graphClosed) { + E.checkNotNull(graphSpace, "graphSpace"); + E.checkNotNull(graph, "graph"); + E.checkNotNull(source, "source"); + E.checkNotNull(store, "store"); + E.checkNotNull(clearCache, "clearCache"); + E.checkNotNull(graphClosed, "graphClosed"); + this.graphSpace = graphSpace; + this.graph = graph; + this.spaceGraphName = graphSpace + "-" + graph; + this.source = source; + this.store = store; + this.clearCache = clearCache; + this.graphClosed = graphClosed; + this.pendingWrite = new AtomicBoolean(false); + } + + private static ScheduledExecutorService newScheduler() { + ScheduledThreadPoolExecutor scheduler = + new ScheduledThreadPoolExecutor(1, r -> { + Thread thread = new Thread(r, "schema-version-reconciler"); + thread.setDaemon(true); + return thread; + }); + scheduler.setRemoveOnCancelPolicy(true); + return scheduler; + } + + static String newVersion(String source) { + return System.currentTimeMillis() + "-" + source + "-" + + VERSION_SEQ.incrementAndGet(); + } + + /** + * Write a new version after a committed schema change. A failed write + * doesn't fail the schema change: it's retried by the next tick, or by + * the next schema change if the reconciler isn't scheduled. + */ + public void bump() { + try { + this.store.write(this.graphSpace, this.graph, + newVersion(this.source)); + } catch (Exception e) { + this.pendingWrite.set(true); + LOG.warn("Schema version write failed for graph '{}', will retry: {}", + this.spaceGraphName, e.toString()); + } + } + + @Override + public void run() { + if (this.stopped || this.graphClosed.getAsBoolean()) { + this.stop(); + return; + } + // An exception escaping here would cancel the scheduled task + try { + // Reset before the write, so a bump() failing meanwhile isn't lost + if (this.pendingWrite.getAndSet(false)) { + try { + this.store.write(this.graphSpace, this.graph, + newVersion(this.source)); + } catch (Exception e) { + this.pendingWrite.set(true); + throw e; + } + LOG.info("Schema version pending write landed for graph '{}'", + this.spaceGraphName); + } + String current = this.store.read(this.graphSpace, this.graph); + if (current == null) { + current = ""; + } + if (this.unreachable) { + this.unreachable = false; + LOG.info("PD reachable again, schema version reconcile " + + "resumed for graph '{}'", this.spaceGraphName); + } + if (current.equals(this.applied)) { + return; + } + // Clear after the read, so the adopted version never claims a + // change that this clear didn't cover + this.clearCache.run(); + String previous = this.applied == null ? NONE : this.applied; + this.applied = current; + LOG.info("Schema cache of graph '{}' cleared by version " + + "reconciler ({} -> {})", this.spaceGraphName, previous, + current.isEmpty() ? NONE : current); + } catch (Exception e) { + if (!this.unreachable) { + this.unreachable = true; + LOG.warn("PD unreachable, schema version reconcile skipped " + + "for graph '{}': {}", this.spaceGraphName, + e.toString()); + } + } + } + + public synchronized void start(long intervalMs) { + E.checkArgument(intervalMs > 0L, + "The reconcile interval must be > 0, but got %s", + intervalMs); + this.stopped = false; + // Random first delay, so servers started together don't tick together + long delay = ThreadLocalRandom.current().nextLong(intervalMs + 1L); + this.future = SCHEDULER.scheduleWithFixedDelay(this, delay, intervalMs, + TimeUnit.MILLISECONDS); + } + + /** + * Restart the task if an Error escaped a tick and cancelled it. Called + * when a schema transaction of the graph is created. + */ + public synchronized void ensureScheduled(long intervalMs) { + if (this.stopped || (this.future != null && !this.future.isDone())) { + return; + } + LOG.warn("Schema version reconciler of graph '{}' is not running, " + + "restarting it", this.spaceGraphName); + this.start(intervalMs); + } + + public synchronized void stop() { + this.stopped = true; + if (this.future != null) { + this.future.cancel(false); + } + } + + public boolean stopped() { + return this.stopped; + } + + public String applied() { + return this.applied; + } + + public boolean pendingWrite() { + return this.pendingWrite.get(); + } + + synchronized ScheduledFuture future() { + return this.future; + } + + public interface VersionStore { + + /** + * @return the stored version, or "" if none was written + */ + String read(String graphSpace, String graph); + + void write(String graphSpace, String graph, String version); + } +} diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/config/CoreOptions.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/config/CoreOptions.java index bbd634eb73..f10bc7406c 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/config/CoreOptions.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/config/CoreOptions.java @@ -529,6 +529,31 @@ public class CoreOptions extends OptionHolder { 10000L ); + public static final ConfigOption SCHEMA_SYNC_ENABLED = + new ConfigOption<>( + "schema.sync.enabled", + "Whether to write a per-graph schema version to PD on every " + + "schema change and check it periodically, so that a schema " + + "change made through another server clears this server's " + + "schema cache. Only for the hstore backend. Set false to " + + "keep the previous behavior: no version writes, no checks.", + disallowEmpty(), + true + ); + + public static final ConfigOption SCHEMA_SYNC_RECONCILE_INTERVAL = + new ConfigOption<>( + "schema.sync.reconcile_interval", + "The interval in seconds to check the schema version of a " + + "graph in PD. A schema change made through another server " + + "is visible on this server within about this interval. " + + "0 means never check: the version is still written for " + + "other servers, and a failed write is retried by the next " + + "schema change.", + rangeInt(0, 3600), + 10 + ); + public static final ConfigOption SCHEMA_INDEX_REBUILD_USING_PUSHDOWN = new ConfigOption<>( "schema.index_rebuild_using_pushdown", diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/MetaManager.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/MetaManager.java index b30b505d32..c5e78c944d 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/MetaManager.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/MetaManager.java @@ -81,6 +81,7 @@ public class MetaManager { public static final String META_PATH_CONF = "CONF"; public static final String META_PATH_GRAPH = "GRAPH"; public static final String META_PATH_SCHEMA = "SCHEMA"; + public static final String META_PATH_SCHEMA_VERSION = "SCHEMA_VERSION"; public static final String META_PATH_PROPERTY_KEY = "PROPERTY_KEY"; public static final String META_PATH_VERTEX_LABEL = "VERTEX_LABEL"; public static final String META_PATH_EDGE_LABEL = "EDGE_LABEL"; @@ -555,6 +556,19 @@ public void notifySchemaCacheClear(String graphSpace, String graph, this.graphMetaManager.notifySchemaCacheClear(graphSpace, graph, source); } + public String getSchemaVersion(String graphSpace, String graph) { + return this.graphMetaManager.getSchemaVersion(graphSpace, graph); + } + + public void putSchemaVersion(String graphSpace, String graph, + String version) { + this.graphMetaManager.putSchemaVersion(graphSpace, graph, version); + } + + public void deleteSchemaVersion(String graphSpace, String graph) { + this.graphMetaManager.deleteSchemaVersion(graphSpace, graph); + } + public void notifyGraphCacheClear(String graphSpace, String graph) { this.graphMetaManager.notifyGraphCacheClear(graphSpace, graph); } diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/GraphMetaManager.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/GraphMetaManager.java index 52b2c3946e..201c1bd6ac 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/GraphMetaManager.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/GraphMetaManager.java @@ -30,6 +30,7 @@ import static org.apache.hugegraph.meta.MetaManager.META_PATH_JOIN; import static org.apache.hugegraph.meta.MetaManager.META_PATH_REMOVE; import static org.apache.hugegraph.meta.MetaManager.META_PATH_SCHEMA; +import static org.apache.hugegraph.meta.MetaManager.META_PATH_SCHEMA_VERSION; import static org.apache.hugegraph.meta.MetaManager.META_PATH_SYS_GRAPH_CONF; import static org.apache.hugegraph.meta.MetaManager.META_PATH_UPDATE; import static org.apache.hugegraph.meta.MetaManager.META_PATH_VERTEX_LABEL; @@ -105,6 +106,25 @@ public void notifySchemaCacheClear(String graphSpace, String graph, graphName(graphSpace, graph), source)); } + /** + * Returns the schema version of the graph, or "" if it was never written. + * The value is opaque and only compared for equality. + */ + public String getSchemaVersion(String graphSpace, String graph) { + String version = this.metaDriver.get( + this.schemaVersionKey(graphSpace, graph)); + return version == null ? "" : version; + } + + public void putSchemaVersion(String graphSpace, String graph, + String version) { + this.metaDriver.put(this.schemaVersionKey(graphSpace, graph), version); + } + + public void deleteSchemaVersion(String graphSpace, String graph) { + this.metaDriver.delete(this.schemaVersionKey(graphSpace, graph)); + } + public void notifyGraphCacheClear(String graphSpace, String graph) { this.metaDriver.put(this.graphCacheClearKey(), graphName(graphSpace, graph)); @@ -273,6 +293,18 @@ private String schemaCacheClearKey() { META_PATH_CLEAR); } + private String schemaVersionKey(String graphSpace, String graph) { + // HUGEGRAPH/{cluster}/SCHEMA_VERSION/{graphspace}/{graph} + // Not under GRAPHSPACE/{graphspace}/{graph}/SCHEMA: clearAllSchema() + // deletes that prefix, which would also match this key. + return String.join(META_PATH_DELIMITER, + META_PATH_HUGEGRAPH, + this.cluster, + META_PATH_SCHEMA_VERSION, + graphSpace, + graph); + } + private String graphCacheClearKey() { // HUGEGRAPH/{cluster}/EVENT/GRAPH/GRAPH/CLEAR return String.join(META_PATH_DELIMITER, diff --git a/hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties b/hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties index 26adfe0183..1421c94911 100644 --- a/hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties +++ b/hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties @@ -4,6 +4,10 @@ gremlin.graph=org.apache.hugegraph.HugeFactory # cache config #schema.cache_capacity=100000 +# hstore only: check the schema version in PD every 10s, so that schema +# changes made through other servers are visible here within about 10s +#schema.sync.enabled=true +#schema.sync.reconcile_interval=10 # vertex-cache default is 1000w, 10min expired vertex.cache_type=l2 #vertex.cache_capacity=10000000 diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java index 0e010ae5f7..e9e300fd54 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java @@ -43,6 +43,7 @@ import org.apache.hugegraph.unit.cache.CachedGraphTransactionTest; import org.apache.hugegraph.unit.cache.CachedSchemaTransactionTest; import org.apache.hugegraph.unit.cache.RamTableTest; +import org.apache.hugegraph.unit.cache.SchemaVersionReconcilerTest; import org.apache.hugegraph.unit.cmd.InitStoreConfigTest; import org.apache.hugegraph.unit.core.AnalyzerTest; import org.apache.hugegraph.unit.core.BackendMutationTest; @@ -57,6 +58,7 @@ import org.apache.hugegraph.unit.core.ExceptionTest; import org.apache.hugegraph.unit.core.GraphManagerAdminInitTest; import org.apache.hugegraph.unit.core.GraphManagerConfigTest; +import org.apache.hugegraph.unit.core.GraphManagerDropGraphTest; import org.apache.hugegraph.unit.core.HstoreSessionsTest; import org.apache.hugegraph.unit.core.IdHolderTest; import org.apache.hugegraph.unit.core.LocksTableTest; @@ -132,6 +134,7 @@ CacheTest.OffheapCacheTest.class, CacheTest.LevelCacheTest.class, CachedSchemaTransactionTest.class, + SchemaVersionReconcilerTest.class, MetaManagerSchemaCacheClearEventTest.class, EtcdMetaDriverTest.class, CachedGraphTransactionTest.class, @@ -172,6 +175,7 @@ ExceptionTest.class, GraphManagerAdminInitTest.class, GraphManagerConfigTest.class, + GraphManagerDropGraphTest.class, HstoreSessionsTest.class, BackendStoreInfoTest.class, TraversalUtilTest.class, diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java index 9e7aec3842..bf67fc2e79 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java @@ -31,6 +31,7 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; +import org.apache.commons.configuration2.PropertiesConfiguration; import org.apache.hugegraph.HugeFactory; import org.apache.hugegraph.HugeGraph; import org.apache.hugegraph.HugeGraphParams; @@ -39,13 +40,17 @@ import org.apache.hugegraph.backend.cache.CacheManager; import org.apache.hugegraph.backend.cache.CachedSchemaTransaction; import org.apache.hugegraph.backend.cache.CachedSchemaTransactionV2; +import org.apache.hugegraph.backend.cache.SchemaVersionReconciler; import org.apache.hugegraph.backend.id.Id; import org.apache.hugegraph.backend.id.IdGenerator; +import org.apache.hugegraph.config.HugeConfig; import org.apache.hugegraph.event.EventHub; import org.apache.hugegraph.event.EventListener; import org.apache.hugegraph.meta.MetaDriver; import org.apache.hugegraph.meta.MetaManager; import org.apache.hugegraph.meta.managers.GraphMetaManager; +import org.apache.hugegraph.meta.managers.SchemaMetaManager; +import org.apache.hugegraph.schema.PropertyKey; import org.apache.hugegraph.schema.SchemaElement; import org.apache.hugegraph.testutil.Assert; import org.apache.hugegraph.testutil.Whitebox; @@ -608,15 +613,276 @@ public void testClearSchemaCacheClearsArrayAttachmentMaps() } } - // TASK_SYNC_DELETION gating of removeSchema notifications and the - // unconditional addSchema notification require an initialised - // CachedSchemaTransactionV2 instance, which in turn needs an hstore - // backend and a connected MetaManager. Both prerequisites are out of - // scope for this unit test class. They are exercised end-to-end by the - // hstore integration tests in CoreTestSuite. TODO(#2617): port these - // assertions into a dedicated CachedSchemaTransactionV2IT once - // mockito-inline becomes available so MetaManager.instance() can be - // stubbed without an hstore cluster. + // The schema version reconciler and the guarded cache updates are covered + // by SchemaVersionReconcilerTest and the V2 tests below. The call sites that + // write the version (add/update/remove/clear), the TASK_SYNC_DELETION + // gating of removeSchema notifications and the unconditional addSchema + // notification need an initialised CachedSchemaTransactionV2, so an + // hstore backend and a connected MetaManager; they are exercised by the + // hstore integration tests in CoreTestSuite. + + @Test + public void testClearV2SchemaCacheBumpsGeneration() { + String graphName = "DEFAULT-generation-v2"; + Cache idCache = v2IdCache(graphName); + Object arrayCaches = idCache.attachment(newV2SchemaCaches(10)); + try { + long before = generation(arrayCaches); + Whitebox.invokeStatic(CachedSchemaTransactionV2.class, + new Class[]{String.class}, + "clearSchemaCache", graphName); + Assert.assertEquals(before + 1L, generation(arrayCaches)); + + // SchemaCaches.clear() alone doesn't change the generation + clearV2SchemaCaches(arrayCaches); + Assert.assertEquals(before + 1L, generation(arrayCaches)); + } finally { + clearV2SchemaCaches(arrayCaches); + idCache.clear(); + } + } + + @Test + public void testV2IdCacheHitIsNotPromotedAcrossRemoteClear() { + Id id = IdGenerator.of(1); + PropertyKey pk = new FakeObjects("unit-test-v2").newPropertyKey(id, + "pk"); + CachedSchemaTransactionV2 tx = v2Tx(); + Object arrayCaches = Whitebox.getInternalState(tx, "arrayCaches"); + Cache idCache = Whitebox.getInternalState(tx, "idCache"); + + // Without a clear the id cache hit is promoted to the array cache + Mockito.when(idCache.get(Mockito.any())).thenReturn(pk); + Assert.assertSame(pk, getV2Schema(tx, id)); + Assert.assertSame(pk, getV2SchemaCache(arrayCaches, + HugeType.PROPERTY_KEY, id)); + + // A remote clear between the id cache read and the promotion + clearV2SchemaCaches(arrayCaches); + Mockito.when(idCache.get(Mockito.any())).thenAnswer(invocation -> { + nextGeneration(arrayCaches); + return pk; + }); + Assert.assertSame(pk, getV2Schema(tx, id)); + Assert.assertNull(getV2SchemaCache(arrayCaches, + HugeType.PROPERTY_KEY, id)); + } + + @Test + public void testV2StorageReadIsNotCachedAcrossRemoteClear() { + Id id = IdGenerator.of(1); + PropertyKey pk = new FakeObjects("unit-test-v2").newPropertyKey(id, + "pk"); + CachedSchemaTransactionV2 tx = v2Tx(); + Object arrayCaches = Whitebox.getInternalState(tx, "arrayCaches"); + Cache idCache = Whitebox.getInternalState(tx, "idCache"); + SchemaMetaManager meta = Whitebox.getInternalState(tx, + "schemaMetaManager"); + + Mockito.when(meta.getPropertyKey(Mockito.any(), Mockito.any(), + Mockito.eq(id))) + .thenAnswer(invocation -> { + nextGeneration(arrayCaches); + return pk; + }); + Assert.assertSame(pk, getV2Schema(tx, id)); + Mockito.verify(idCache, Mockito.never()).update(Mockito.any(), + Mockito.any()); + Assert.assertNull(getV2SchemaCache(arrayCaches, + HugeType.PROPERTY_KEY, id)); + + // Without a clear the storage read is cached + Mockito.when(meta.getPropertyKey(Mockito.any(), Mockito.any(), + Mockito.eq(id))) + .thenReturn(pk); + Assert.assertSame(pk, getV2Schema(tx, id)); + Mockito.verify(idCache).update(Mockito.any(), Mockito.eq(pk)); + Assert.assertSame(pk, getV2SchemaCache(arrayCaches, + HugeType.PROPERTY_KEY, id)); + } + + @Test + public void testV2AllSchemaIsNotCachedAcrossRemoteClear() + throws Exception { + PropertyKey pk = new FakeObjects("unit-test-v2").newPropertyKey( + IdGenerator.of(1), "pk"); + CachedSchemaTransactionV2 tx = v2Tx(); + Object arrayCaches = Whitebox.getInternalState(tx, "arrayCaches"); + Cache idCache = Whitebox.getInternalState(tx, "idCache"); + SchemaMetaManager meta = Whitebox.getInternalState(tx, + "schemaMetaManager"); + Map cachedTypes = readField(arrayCaches, + "cachedTypes"); + + Mockito.when(meta.getPropertyKeys(Mockito.any(), Mockito.any())) + .thenAnswer(invocation -> { + nextGeneration(arrayCaches); + return Collections.singletonList(pk); + }); + Assert.assertEquals(1, getV2AllSchema(tx).size()); + Mockito.verify(idCache, Mockito.never()).update(Mockito.any(), + Mockito.any()); + Assert.assertNull(cachedTypes.get(HugeType.PROPERTY_KEY)); + + Mockito.when(meta.getPropertyKeys(Mockito.any(), Mockito.any())) + .thenReturn(Collections.singletonList(pk)); + Assert.assertEquals(1, getV2AllSchema(tx).size()); + Mockito.verify(idCache).update(Mockito.any(), Mockito.eq(pk)); + Assert.assertEquals(true, cachedTypes.get(HugeType.PROPERTY_KEY)); + } + + @Test + public void testV2LocalClearKeepsGeneration() { + CachedSchemaTransactionV2 tx = v2Tx(); + Object arrayCaches = Whitebox.getInternalState(tx, "arrayCaches"); + long before = generation(arrayCaches); + + // Only a clear caused by another server changes the generation, so a + // local clear (e.g. the name miss reload) can't drop other readers + tx.clearCache(false); + Assert.assertEquals(before, generation(arrayCaches)); + } + + @Test + public void testV2ReconcilerStartsOnceAndStopsOnGraphClose() { + HugeGraphParams params = mockV2Params("DEFAULT-registry-v2", + syncConfig(true, 3600)); + try { + SchemaVersionReconciler first = ensureReconciler(params); + Assert.assertNotNull(first); + Assert.assertSame(first, ensureReconciler(params)); + + CachedSchemaTransactionV2.stopReconciler("DEFAULT-registry-v2"); + Assert.assertTrue(first.stopped()); + + SchemaVersionReconciler second = ensureReconciler(params); + Assert.assertNotSame(first, second); + Assert.assertFalse(second.stopped()); + + // A closed graph gets no new reconciler + CachedSchemaTransactionV2.stopReconciler("DEFAULT-registry-v2"); + Mockito.when(params.closed()).thenReturn(true); + Assert.assertNull(ensureReconciler(params)); + } finally { + CachedSchemaTransactionV2.stopReconciler("DEFAULT-registry-v2"); + } + } + + @Test + public void testV2ReconcilerHonoursSyncOptions() throws Exception { + HugeGraphParams disabled = mockV2Params("DEFAULT-disabled-v2", + syncConfig(false, 10)); + Assert.assertNull(ensureReconciler(disabled)); + + MetaDriver mockDriver = Mockito.mock(MetaDriver.class); + Object previous = swapGraphMetaManager( + new GraphMetaManager(mockDriver, "c")); + HugeGraphParams noPolling = mockV2Params("DEFAULT-nopoll-v2", + syncConfig(true, 0)); + try { + SchemaVersionReconciler reconciler = ensureReconciler(noPolling); + Assert.assertNotNull(reconciler); + Assert.assertNull(Whitebox.invoke(SchemaVersionReconciler.class, + "future", reconciler)); + + // The version is still written for the other servers + reconciler.bump(); + Mockito.verify(mockDriver).put( + Mockito.eq("HUGEGRAPH/c/SCHEMA_VERSION/DEFAULT/nopoll-v2"), + Mockito.anyString()); + Assert.assertFalse(reconciler.pendingWrite()); + } finally { + CachedSchemaTransactionV2.stopReconciler("DEFAULT-nopoll-v2"); + swapGraphMetaManager(previous); + } + } + + @Test + public void testSchemaVersionKeyAndValue() { + MetaDriver mockDriver = Mockito.mock(MetaDriver.class); + GraphMetaManager manager = new GraphMetaManager(mockDriver, "c"); + String key = "HUGEGRAPH/c/SCHEMA_VERSION/DEFAULT/g"; + + manager.putSchemaVersion("DEFAULT", "g", "v1"); + Mockito.verify(mockDriver).put(key, "v1"); + + Mockito.when(mockDriver.get(key)).thenReturn(null); + Assert.assertEquals("", manager.getSchemaVersion("DEFAULT", "g")); + Mockito.when(mockDriver.get(key)).thenReturn("v2"); + Assert.assertEquals("v2", manager.getSchemaVersion("DEFAULT", "g")); + + manager.deleteSchemaVersion("DEFAULT", "g"); + Mockito.verify(mockDriver).delete(key); + } + + private static long generation(Object arrayCaches) { + Long generation = Whitebox.invoke(arrayCaches.getClass(), + "generation", arrayCaches); + return generation; + } + + private static void nextGeneration(Object arrayCaches) { + Whitebox.invoke(arrayCaches.getClass(), "nextGeneration", arrayCaches); + } + + @SuppressWarnings("unchecked") + private static CachedSchemaTransactionV2 v2Tx() { + // The constructor needs an hstore backend: build the instance without + // it and set only the fields the read paths use + CachedSchemaTransactionV2 tx = Mockito.mock( + CachedSchemaTransactionV2.class, + Mockito.withSettings().defaultAnswer(Mockito.CALLS_REAL_METHODS)); + Cache idCache = Mockito.mock(Cache.class); + Mockito.when(idCache.capacity()).thenReturn(100L); + Whitebox.setInternalState(tx, "idCache", idCache); + Whitebox.setInternalState(tx, "nameCache", Mockito.mock(Cache.class)); + Whitebox.setInternalState(tx, "arrayCaches", newV2SchemaCaches(10)); + Whitebox.setInternalState(tx, "schemaMetaManager", + Mockito.mock(SchemaMetaManager.class)); + return tx; + } + + private static SchemaElement getV2Schema(CachedSchemaTransactionV2 tx, + Id id) { + return Whitebox.invoke(CachedSchemaTransactionV2.class, + new Class[]{HugeType.class, Id.class}, + "getSchema", tx, HugeType.PROPERTY_KEY, id); + } + + private static List getV2AllSchema( + CachedSchemaTransactionV2 tx) { + return Whitebox.invoke(CachedSchemaTransactionV2.class, + new Class[]{HugeType.class}, + "getAllSchema", tx, HugeType.PROPERTY_KEY); + } + + private static HugeConfig syncConfig(boolean enabled, int interval) { + PropertiesConfiguration conf = new PropertiesConfiguration(); + conf.setProperty("schema.sync.enabled", enabled); + conf.setProperty("schema.sync.reconcile_interval", interval); + return new HugeConfig(conf); + } + + private static HugeGraphParams mockV2Params(String spaceGraphName, + HugeConfig config) { + String[] parts = spaceGraphName.split("-", 2); + HugeGraph graph = Mockito.mock(HugeGraph.class); + Mockito.when(graph.spaceGraphName()).thenReturn(spaceGraphName); + Mockito.when(graph.graphSpace()).thenReturn(parts[0]); + Mockito.when(graph.name()).thenReturn(parts[1]); + HugeGraphParams params = Mockito.mock(HugeGraphParams.class); + Mockito.when(params.graph()).thenReturn(graph); + Mockito.when(params.configuration()).thenReturn(config); + Mockito.when(params.closed()).thenReturn(false); + return params; + } + + private static SchemaVersionReconciler ensureReconciler( + HugeGraphParams params) { + return Whitebox.invokeStatic(CachedSchemaTransactionV2.class, + new Class[]{HugeGraphParams.class}, + "ensureReconciler", params); + } @Test public void testHandleSchemaCacheClearEventSkipsLocalSource() diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/SchemaVersionReconcilerTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/SchemaVersionReconcilerTest.java new file mode 100644 index 0000000000..049c8cbe73 --- /dev/null +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/SchemaVersionReconcilerTest.java @@ -0,0 +1,355 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with this + * work for additional information regarding copyright ownership. The ASF + * licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hugegraph.unit.cache; + +import java.util.HashSet; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.hugegraph.HugeException; +import org.apache.hugegraph.backend.cache.SchemaVersionReconciler; +import org.apache.hugegraph.backend.cache.SchemaVersionReconciler.VersionStore; +import org.apache.hugegraph.testutil.Assert; +import org.apache.hugegraph.testutil.Whitebox; +import org.apache.hugegraph.unit.BaseUnitTest; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +public class SchemaVersionReconcilerTest extends BaseUnitTest { + + private static final String SPACE = "DEFAULT"; + private static final String GRAPH = "g"; + + private MemoryStore store; + private AtomicInteger clears; + private AtomicBoolean closed; + private SchemaVersionReconciler reconciler; + + @Before + public void setup() { + this.store = new MemoryStore(); + this.clears = new AtomicInteger(); + this.closed = new AtomicBoolean(false); + this.reconciler = this.newReconciler(this.store); + } + + @After + public void teardown() { + this.reconciler.stop(); + } + + private SchemaVersionReconciler newReconciler(VersionStore store) { + return new SchemaVersionReconciler(SPACE, GRAPH, "src", store, + this.clears::incrementAndGet, + this.closed::get); + } + + @Test + public void testFirstTickClearsAndAdoptsAbsentVersion() { + this.reconciler.run(); + Assert.assertEquals(1, this.clears.get()); + Assert.assertEquals("", this.reconciler.applied()); + } + + @Test + public void testTickDoesNothingWhenVersionUnchanged() { + this.store.put("v1"); + this.reconciler.run(); + this.reconciler.run(); + this.reconciler.run(); + Assert.assertEquals(1, this.clears.get()); + Assert.assertEquals("v1", this.reconciler.applied()); + } + + @Test + public void testTickClearsOnChangeFromAnotherServer() { + this.reconciler.run(); + this.store.put("v2"); + this.reconciler.run(); + Assert.assertEquals(2, this.clears.get()); + Assert.assertEquals("v2", this.reconciler.applied()); + } + + @Test + public void testManyChangesBetweenTicksClearOnce() { + this.reconciler.run(); + for (int i = 0; i < 1000; i++) { + this.store.put("v" + i); + } + this.reconciler.run(); + Assert.assertEquals(2, this.clears.get()); + Assert.assertEquals("v999", this.reconciler.applied()); + } + + @Test + public void testNewVersionIsUnique() { + Set versions = new HashSet<>(); + for (int i = 0; i < 1000; i++) { + String version = Whitebox.invokeStatic( + SchemaVersionReconciler.class, + new Class[]{String.class}, "newVersion", "s"); + Assert.assertContains("-s-", version); + versions.add(version); + } + Assert.assertEquals(1000, versions.size()); + } + + @Test + public void testBumpWritesVersionAndOwnChangeIsNotSkipped() { + this.reconciler.run(); + this.reconciler.bump(); + String written = this.store.get(); + Assert.assertFalse(written.isEmpty()); + // bump() never adopts its own version + Assert.assertEquals("", this.reconciler.applied()); + + this.reconciler.run(); + Assert.assertEquals(2, this.clears.get()); + Assert.assertEquals(written, this.reconciler.applied()); + } + + @Test + public void testBumpFailureIsRetriedByNextTick() { + this.store.failWrites(1, new HugeException("pd down")); + this.reconciler.bump(); + Assert.assertTrue(this.reconciler.pendingWrite()); + Assert.assertEquals("", this.store.get()); + + this.reconciler.run(); + Assert.assertFalse(this.reconciler.pendingWrite()); + Assert.assertFalse(this.store.get().isEmpty()); + Assert.assertEquals(this.store.get(), this.reconciler.applied()); + } + + @Test + public void testBumpFailureWithNullPointerIsRetried() { + // PdMetaDriver.put() fails this way when every PD is unreachable + this.store.failWrites(1, new NullPointerException()); + this.reconciler.bump(); + Assert.assertTrue(this.reconciler.pendingWrite()); + this.reconciler.run(); + Assert.assertFalse(this.reconciler.pendingWrite()); + } + + @Test + public void testBumpFailureDuringRetryIsNotLost() { + this.store.failWrites(1, new HugeException("pd down")); + this.reconciler.bump(); + // The retry of the tick writes, then a new change fails to write + this.store.onWrite(() -> { + this.store.failWrites(1, new HugeException("pd down")); + this.reconciler.bump(); + }); + this.reconciler.run(); + Assert.assertTrue(this.reconciler.pendingWrite()); + } + + @Test + public void testBumpDoesNotSwallowErrors() { + this.store.failWrites(1, new AssertionError("fatal")); + Assert.assertThrows(AssertionError.class, () -> { + this.reconciler.bump(); + }); + Assert.assertFalse(this.reconciler.pendingWrite()); + } + + @Test + public void testTickSurvivesReadFailures() { + this.store.put("v1"); + this.reconciler.run(); + + this.store.put("v2"); + this.store.failReads(2); + this.reconciler.run(); + this.reconciler.run(); + Assert.assertEquals(1, this.clears.get()); + Assert.assertEquals("v1", this.reconciler.applied()); + + this.reconciler.run(); + Assert.assertEquals(2, this.clears.get()); + Assert.assertEquals("v2", this.reconciler.applied()); + } + + @Test + public void testAdoptsVersionReadBeforeClear() { + // A change landing while the cache is being cleared must not be + // marked as covered by that clear + AtomicReference next = new AtomicReference<>("v2"); + SchemaVersionReconciler reconciler = new SchemaVersionReconciler( + SPACE, GRAPH, "src", this.store, () -> { + this.clears.incrementAndGet(); + String value = next.getAndSet(null); + if (value != null) { + this.store.put(value); + } + }, this.closed::get); + this.store.put("v1"); + + reconciler.run(); + Assert.assertEquals("v1", reconciler.applied()); + reconciler.run(); + Assert.assertEquals(2, this.clears.get()); + Assert.assertEquals("v2", reconciler.applied()); + } + + @Test + public void testClosedGraphStopsReconciler() { + this.closed.set(true); + this.reconciler.run(); + Assert.assertTrue(this.reconciler.stopped()); + Assert.assertEquals(0, this.clears.get()); + Assert.assertEquals(0, this.store.reads.get()); + } + + @Test + public void testScheduledTicksStopAfterStop() throws Exception { + this.reconciler.start(20L); + waitFor(() -> this.store.reads.get() >= 2); + Assert.assertTrue(this.store.readThread.get().isDaemon()); + Assert.assertEquals("schema-version-reconciler", + this.store.readThread.get().getName()); + + this.reconciler.stop(); + Thread.sleep(50L); + int reads = this.store.reads.get(); + Thread.sleep(100L); + Assert.assertEquals(reads, this.store.reads.get()); + } + + @Test + public void testEnsureScheduledRestartsTaskKilledByError() + throws Exception { + this.store.failReads(1, new AssertionError("fatal")); + this.reconciler.start(20L); + ScheduledFuture first = Whitebox.invoke( + SchemaVersionReconciler.class, "future", + this.reconciler); + waitFor(first::isDone); + + this.reconciler.ensureScheduled(20L); + ScheduledFuture second = Whitebox.invoke( + SchemaVersionReconciler.class, "future", + this.reconciler); + Assert.assertNotSame(first, second); + waitFor(() -> this.clears.get() >= 1); + this.reconciler.stop(); + } + + @Test + public void testEnsureScheduledKeepsRunningOrStoppedTask() { + this.reconciler.start(60_000L); + ScheduledFuture first = Whitebox.invoke( + SchemaVersionReconciler.class, "future", + this.reconciler); + this.reconciler.ensureScheduled(60_000L); + Assert.assertSame(first, Whitebox.invoke(SchemaVersionReconciler.class, + "future", this.reconciler)); + + this.reconciler.stop(); + this.reconciler.ensureScheduled(60_000L); + Assert.assertSame(first, Whitebox.invoke(SchemaVersionReconciler.class, + "future", this.reconciler)); + Assert.assertTrue(first.isCancelled()); + } + + private static void waitFor(java.util.function.BooleanSupplier condition) + throws InterruptedException { + long deadline = System.currentTimeMillis() + 5000L; + while (!condition.getAsBoolean()) { + if (System.currentTimeMillis() > deadline) { + Assert.fail("Timed out waiting for the condition"); + } + Thread.sleep(5L); + } + } + + private static class MemoryStore implements VersionStore { + + private final Map versions = new ConcurrentHashMap<>(); + private final AtomicInteger reads = new AtomicInteger(); + private final AtomicReference readThread = + new AtomicReference<>(); + private final AtomicInteger readFailures = new AtomicInteger(); + private final AtomicInteger writeFailures = new AtomicInteger(); + private volatile Throwable readFailure; + private volatile Throwable writeFailure; + private volatile Runnable onWrite; + + void put(String version) { + this.versions.put(SPACE + "/" + GRAPH, version); + } + + String get() { + return this.versions.getOrDefault(SPACE + "/" + GRAPH, ""); + } + + void failReads(int times) { + this.failReads(times, new HugeException("pd down")); + } + + void failReads(int times, Throwable failure) { + this.readFailure = failure; + this.readFailures.set(times); + } + + void failWrites(int times, Throwable failure) { + this.writeFailure = failure; + this.writeFailures.set(times); + } + + void onWrite(Runnable action) { + this.onWrite = action; + } + + @Override + public String read(String graphSpace, String graph) { + this.reads.incrementAndGet(); + this.readThread.set(Thread.currentThread()); + if (this.readFailures.getAndDecrement() > 0) { + throwUnchecked(this.readFailure); + } + return this.versions.getOrDefault(graphSpace + "/" + graph, ""); + } + + @Override + public void write(String graphSpace, String graph, String version) { + if (this.writeFailures.getAndDecrement() > 0) { + throwUnchecked(this.writeFailure); + } + this.versions.put(graphSpace + "/" + graph, version); + Runnable action = this.onWrite; + if (action != null) { + this.onWrite = null; + action.run(); + } + } + + private static void throwUnchecked(Throwable failure) { + if (failure instanceof Error) { + throw (Error) failure; + } + throw (RuntimeException) failure; + } + } +} diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/GraphManagerDropGraphTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/GraphManagerDropGraphTest.java new file mode 100644 index 0000000000..897f0775b5 --- /dev/null +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/GraphManagerDropGraphTest.java @@ -0,0 +1,124 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with this + * work for additional information regarding copyright ownership. The ASF + * licenses this file to You under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hugegraph.unit.core; + +import java.lang.reflect.Field; +import java.util.Map; + +import org.apache.commons.configuration2.PropertiesConfiguration; +import org.apache.hugegraph.HugeException; +import org.apache.hugegraph.HugeGraph; +import org.apache.hugegraph.config.HugeConfig; +import org.apache.hugegraph.core.GraphManager; +import org.apache.hugegraph.event.EventHub; +import org.apache.hugegraph.meta.MetaDriver; +import org.apache.hugegraph.meta.MetaManager; +import org.apache.hugegraph.meta.managers.GraphMetaManager; +import org.apache.hugegraph.meta.managers.SpaceMetaManager; +import org.apache.hugegraph.space.GraphSpace; +import org.apache.hugegraph.task.TaskScheduler; +import org.apache.hugegraph.testutil.Assert; +import org.apache.hugegraph.testutil.Whitebox; +import org.apache.hugegraph.unit.BaseUnitTest; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.mockito.InOrder; +import org.mockito.Mockito; + +public class GraphManagerDropGraphTest extends BaseUnitTest { + + private static final String CLUSTER = "drop-test"; + private static final String VERSION_KEY = + "HUGEGRAPH/drop-test/SCHEMA_VERSION/DEFAULT/g"; + + private MetaDriver driver; + private HugeGraph graph; + private Object originalGraphManager; + private Object originalSpaceManager; + private GraphManager manager; + + @Before + public void setup() throws Exception { + this.driver = Mockito.mock(MetaDriver.class); + this.originalGraphManager = swapMetaManagerField( + "graphMetaManager", new GraphMetaManager(this.driver, CLUSTER)); + this.originalSpaceManager = swapMetaManagerField( + "spaceMetaManager", new SpaceMetaManager(this.driver, CLUSTER)); + this.manager = new GraphManager( + new HugeConfig(new PropertiesConfiguration()), + new EventHub("drop-graph-test")); + // Take the PD branch of dropGraph() with one open graph "DEFAULT-g" + Whitebox.setInternalState(this.manager, "PDExist", true); + this.graph = Mockito.mock(HugeGraph.class); + Mockito.when(this.graph.taskScheduler()) + .thenReturn(Mockito.mock(TaskScheduler.class)); + Map graphs = Whitebox.getInternalState(this.manager, + "graphs"); + graphs.put("DEFAULT-g", this.graph); + Map spaces = Whitebox.getInternalState( + this.manager, "graphSpaces"); + spaces.put("DEFAULT", new GraphSpace("DEFAULT")); + } + + @After + public void teardown() throws Exception { + try { + Whitebox.setInternalState(this.manager, "PDExist", false); + this.manager.close(); + } finally { + swapMetaManagerField("graphMetaManager", this.originalGraphManager); + swapMetaManagerField("spaceMetaManager", this.originalSpaceManager); + } + } + + @Test + public void testDropGraphDeletesSchemaVersionAfterClear() + throws Exception { + this.manager.dropGraph("DEFAULT", "g", true); + + InOrder order = Mockito.inOrder(this.graph, this.driver); + order.verify(this.graph).clearBackend(); + order.verify(this.driver).delete(VERSION_KEY); + order.verify(this.graph, Mockito.atLeastOnce()).close(); + } + + @Test + public void testDropGraphContinuesWhenSchemaVersionDeleteFails() + throws Exception { + Mockito.doThrow(new HugeException("pd down")) + .when(this.driver).delete(VERSION_KEY); + + this.manager.dropGraph("DEFAULT", "g", true); + + Mockito.verify(this.graph, Mockito.atLeastOnce()).close(); + Map graphs = Whitebox.getInternalState(this.manager, + "graphs"); + Assert.assertFalse(graphs.containsKey("DEFAULT-g")); + } + + private static Object swapMetaManagerField(String field, + Object replacement) + throws Exception { + Field f = MetaManager.class.getDeclaredField(field); + f.setAccessible(true); + Object previous = f.get(MetaManager.instance()); + f.set(MetaManager.instance(), replacement); + return previous; + } +} From be9324233c4871fe980141639e9946e7b3e30529 Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Thu, 24 Sep 2026 19:10:06 +0530 Subject: [PATCH 2/9] doc(server): state the rollout conditions of schema sync The convergence bound needs every Server of the graph on a release with these options, sync enabled, polling on and PD reachable. Say so in the conf template, with what to do after a rolling upgrade: older Servers write no version for schema updates. --- .../src/assembly/static/conf/graphs/hugegraph.properties | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties b/hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties index 1421c94911..85d334c8ac 100644 --- a/hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties +++ b/hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties @@ -4,8 +4,13 @@ gremlin.graph=org.apache.hugegraph.HugeFactory # cache config #schema.cache_capacity=100000 -# hstore only: check the schema version in PD every 10s, so that schema -# changes made through other servers are visible here within about 10s +# hstore only: every schema change writes a version to PD and each server +# checks it every reconcile_interval seconds, so a change made through another +# server is visible here within about that time. This holds only when every +# server of the graph runs a release with these options, with sync enabled, +# the interval above 0 and PD reachable. Older servers write no version for +# schema updates, so after a rolling upgrade make one schema change, or +# restart the servers, to reload every server's schema cache. #schema.sync.enabled=true #schema.sync.reconcile_interval=10 # vertex-cache default is 1000w, 10min expired From 070a96186d1537c1e3e99973ee52fe4499e8ee75 Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Fri, 25 Sep 2026 01:48:15 +0530 Subject: [PATCH 3/9] fix(server): guard the local schema write cache update addSchema() and updateSchema() cached the written element without the generation check the read paths use. A remote clear landing between the storage write and the cache update could then be undone by an element another server had already replaced. Capture the generation before the write and update the caches under the arrayCaches lock only if it is unchanged. --- .../cache/CachedSchemaTransactionV2.java | 36 ++++++++++++------- .../cache/CachedSchemaTransactionTest.java | 32 +++++++++++++++++ 2 files changed, 56 insertions(+), 12 deletions(-) diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java index 485368c952..53d1f72346 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java @@ -397,9 +397,10 @@ private void invalidateCache(HugeType type, Id id) { @Override protected void updateSchema(SchemaElement schema, Consumer updateCallback) { + long generation = this.arrayCaches.generation(); super.updateSchema(schema, updateCallback); - this.updateCache(schema); + this.updateCache(schema, generation); // No meta event here: one per updateSchemaStatus() call from background // jobs would be a broadcast storm. The schema version has no fan-out, // other servers check it at most once per reconcile interval. @@ -408,9 +409,10 @@ protected void updateSchema(SchemaElement schema, @Override protected void addSchema(SchemaElement schema) { + long generation = this.arrayCaches.generation(); super.addSchema(schema); - this.updateCache(schema); + this.updateCache(schema, generation); this.bumpSchemaVersion(); // Schema additions must always propagate to remote nodes regardless @@ -418,19 +420,29 @@ protected void addSchema(SchemaElement schema) { this.notifySchemaCacheClear(); } - private void updateCache(SchemaElement schema) { - this.resetCachedAllIfReachedCapacity(); + private void updateCache(SchemaElement schema, long generation) { + synchronized (this.arrayCaches) { + /* + * Skip the update if a remote change cleared the caches after + * the storage write: another server may have written a newer + * element meanwhile, and the next read loads it from storage. + */ + if (generation != this.arrayCaches.generation()) { + return; + } + this.resetCachedAllIfReachedCapacity(); - // update id cache - Id prefixedId = generateId(schema.type(), schema.id()); - this.idCache.update(prefixedId, schema); + // update id cache + Id prefixedId = generateId(schema.type(), schema.id()); + this.idCache.update(prefixedId, schema); - // update name cache - Id prefixedName = generateId(schema.type(), schema.name()); - this.nameCache.update(prefixedName, schema); + // update name cache + Id prefixedName = generateId(schema.type(), schema.name()); + this.nameCache.update(prefixedName, schema); - // update optimized array cache - this.arrayCaches.updateIfNeeded(schema); + // update optimized array cache + this.arrayCaches.updateIfNeeded(schema); + } } @Override diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java index bf67fc2e79..bcb1f51947 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java @@ -731,6 +731,31 @@ public void testV2AllSchemaIsNotCachedAcrossRemoteClear() Assert.assertEquals(true, cachedTypes.get(HugeType.PROPERTY_KEY)); } + @Test + public void testV2LocalWriteIsNotCachedAcrossRemoteClear() { + Id id = IdGenerator.of(1); + PropertyKey pk = new FakeObjects("unit-test-v2").newPropertyKey(id, + "pk"); + CachedSchemaTransactionV2 tx = v2Tx(); + Object arrayCaches = Whitebox.getInternalState(tx, "arrayCaches"); + Cache idCache = Whitebox.getInternalState(tx, "idCache"); + + // A remote clear between the storage write and the cache update + long generation = generation(arrayCaches); + nextGeneration(arrayCaches); + updateV2Cache(tx, pk, generation); + Mockito.verify(idCache, Mockito.never()).update(Mockito.any(), + Mockito.any()); + Assert.assertNull(getV2SchemaCache(arrayCaches, + HugeType.PROPERTY_KEY, id)); + + // Without a clear the written element is cached + updateV2Cache(tx, pk, generation(arrayCaches)); + Mockito.verify(idCache).update(Mockito.any(), Mockito.eq(pk)); + Assert.assertSame(pk, getV2SchemaCache(arrayCaches, + HugeType.PROPERTY_KEY, id)); + } + @Test public void testV2LocalClearKeepsGeneration() { CachedSchemaTransactionV2 tx = v2Tx(); @@ -849,6 +874,13 @@ private static SchemaElement getV2Schema(CachedSchemaTransactionV2 tx, "getSchema", tx, HugeType.PROPERTY_KEY, id); } + private static void updateV2Cache(CachedSchemaTransactionV2 tx, + SchemaElement schema, long generation) { + Whitebox.invoke(CachedSchemaTransactionV2.class, + new Class[]{SchemaElement.class, long.class}, + "updateCache", tx, schema, generation); + } + private static List getV2AllSchema( CachedSchemaTransactionV2 tx) { return Whitebox.invoke(CachedSchemaTransactionV2.class, From 7f5672a7c8fc5b7066ea631432bfaf73bc95b60d Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Fri, 25 Sep 2026 02:51:51 +0530 Subject: [PATCH 4/9] fix(server): drop the written element when a clear races the write 070a9618 skipped the cache update of a local write if the generation changed during it. The reconciler also clears on this server's own version writes, so a reader could cache the old element after that clear and the skipped update left it there: an index label stayed CREATING in the cache after its rebuild set CREATED, failing two hstore core tests. Invalidate the element and reset the cached-all flag of its type instead, so the next read loads it from storage. --- .../cache/CachedSchemaTransactionV2.java | 19 +++++++----- .../cache/CachedSchemaTransactionTest.java | 29 ++++++++++++------- 2 files changed, 31 insertions(+), 17 deletions(-) diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java index 53d1f72346..1f38b84ba6 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java @@ -421,23 +421,28 @@ protected void addSchema(SchemaElement schema) { } private void updateCache(SchemaElement schema, long generation) { + Id prefixedId = generateId(schema.type(), schema.id()); + Id prefixedName = generateId(schema.type(), schema.name()); synchronized (this.arrayCaches) { - /* - * Skip the update if a remote change cleared the caches after - * the storage write: another server may have written a newer - * element meanwhile, and the next read loads it from storage. - */ if (generation != this.arrayCaches.generation()) { + /* + * A remote change cleared the caches during this write. Since + * then a reader may have cached an older copy, and another + * server may have written a newer one, so drop the element + * and let the next read load it from storage. + */ + this.idCache.invalidate(prefixedId); + this.nameCache.invalidate(prefixedName); + this.arrayCaches.remove(schema.type(), schema.id()); + this.resetCachedAll(schema.type()); return; } this.resetCachedAllIfReachedCapacity(); // update id cache - Id prefixedId = generateId(schema.type(), schema.id()); this.idCache.update(prefixedId, schema); // update name cache - Id prefixedName = generateId(schema.type(), schema.name()); this.nameCache.update(prefixedName, schema); // update optimized array cache diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java index bcb1f51947..998495ff89 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java @@ -732,28 +732,37 @@ public void testV2AllSchemaIsNotCachedAcrossRemoteClear() } @Test - public void testV2LocalWriteIsNotCachedAcrossRemoteClear() { + public void testV2LocalWriteIsDroppedAcrossRemoteClear() + throws Exception { Id id = IdGenerator.of(1); PropertyKey pk = new FakeObjects("unit-test-v2").newPropertyKey(id, "pk"); CachedSchemaTransactionV2 tx = v2Tx(); Object arrayCaches = Whitebox.getInternalState(tx, "arrayCaches"); Cache idCache = Whitebox.getInternalState(tx, "idCache"); - - // A remote clear between the storage write and the cache update - long generation = generation(arrayCaches); - nextGeneration(arrayCaches); - updateV2Cache(tx, pk, generation); - Mockito.verify(idCache, Mockito.never()).update(Mockito.any(), - Mockito.any()); - Assert.assertNull(getV2SchemaCache(arrayCaches, - HugeType.PROPERTY_KEY, id)); + Map cachedTypes = readField(arrayCaches, + "cachedTypes"); // Without a clear the written element is cached updateV2Cache(tx, pk, generation(arrayCaches)); Mockito.verify(idCache).update(Mockito.any(), Mockito.eq(pk)); Assert.assertSame(pk, getV2SchemaCache(arrayCaches, HugeType.PROPERTY_KEY, id)); + + /* + * A remote clear between the start of the write and the cache update: + * a reader may have cached an older copy since, so the element is + * dropped instead of skipped or overwritten + */ + cachedTypes.put(HugeType.PROPERTY_KEY, true); + long generation = generation(arrayCaches); + nextGeneration(arrayCaches); + updateV2Cache(tx, pk, generation); + Mockito.verify(idCache).update(Mockito.any(), Mockito.eq(pk)); + Mockito.verify(idCache).invalidate(Mockito.any()); + Assert.assertNull(getV2SchemaCache(arrayCaches, + HugeType.PROPERTY_KEY, id)); + Assert.assertEquals(false, cachedTypes.get(HugeType.PROPERTY_KEY)); } @Test From be867b28ce58ff1c71df82b19d8b7c7b49117583 Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Sun, 27 Sep 2026 18:07:17 +0530 Subject: [PATCH 5/9] fix(server): keep the schema id counters on hstore truncate HstoreSession.truncate() called resetIdByKey() for the store's graph name. For the schema store that is {graphspace}/{graph}/m, so PD dropped every schema id counter of the graph by prefix, while the schema itself stays in PD meta and survives the truncate. The truncating Server kept handing out ids from its cached range, but any other Server, or the same one after a restart, got ids from 1 again and reused the ids of existing schema elements. Stop resetting the counters on truncate. --- .../store/hstore/HstoreSessionsImpl.java | 4 +- .../hugegraph/core/MultiGraphsTest.java | 38 +++++++++++++++++++ 2 files changed, 40 insertions(+), 2 deletions(-) diff --git a/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreSessionsImpl.java b/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreSessionsImpl.java index d172651393..d9b2130602 100755 --- a/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreSessionsImpl.java +++ b/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreSessionsImpl.java @@ -801,9 +801,9 @@ public void setMode(GraphMode mode) { @Override public void truncate() throws Exception { + // The schema lives in PD meta and survives a truncate, so the schema + // id counters in PD must survive too: a reset hands out ids again this.graph.truncate(); - HstoreSessionsImpl.getDefaultPdClient() - .resetIdByKey(this.getGraphName()); } @Override diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/MultiGraphsTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/MultiGraphsTest.java index 7ed172acd9..ac7f26968c 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/MultiGraphsTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/core/MultiGraphsTest.java @@ -21,6 +21,7 @@ import java.util.ArrayList; import java.util.Iterator; import java.util.List; +import java.util.Map; import java.util.Objects; import org.apache.commons.configuration2.BaseConfiguration; @@ -31,6 +32,7 @@ import org.apache.hugegraph.backend.id.IdGenerator; import org.apache.hugegraph.backend.store.BackendStoreInfo; import org.apache.hugegraph.backend.store.rocksdb.RocksDBOptions; +import org.apache.hugegraph.backend.tx.IdCounter; import org.apache.hugegraph.config.CoreOptions; import org.apache.hugegraph.exception.ExistedException; import org.apache.hugegraph.masterelection.GlobalMasterInfo; @@ -41,6 +43,7 @@ import org.apache.hugegraph.schema.VertexLabel; import org.apache.hugegraph.testutil.Assert; import org.apache.hugegraph.testutil.Utils; +import org.apache.hugegraph.testutil.Whitebox; import org.apache.tinkerpop.gremlin.structure.T; import org.apache.tinkerpop.gremlin.structure.Vertex; import org.apache.tinkerpop.gremlin.structure.util.GraphFactory; @@ -121,6 +124,41 @@ public void testTruncateBackendKeepsVersionAndResetsSchemaIds() { } } + @Test + public void testHstoreTruncateBackendKeepsSchemaIdCounters() { + Assume.assumeTrue("only hstore keeps the schema id counters in PD", + "hstore".equals(graph().backend())); + + HugeGraph graph = openGraphs("truncate_hs").get(0); + try { + graph.clearBackend(); + graph.initBackend(); + graph.serverStarted(GlobalMasterInfo.master("server-truncate")); + + SchemaManager schema = graph.schema(); + PropertyKey name = schema.propertyKey("name").asText().create(); + + graph.truncateBackend(); + + // The schema lives in PD meta and survives the truncate + Assert.assertEquals(name.id(), schema.getPropertyKey("name").id()); + + // Another Server holds no cached id range, it asks PD for one: drop + // the ranges this process cached for the graph's schema counters + Map ranges = Whitebox.getInternalState(IdCounter.class, "ids"); + String prefix = String.join("/", graph.graphSpace(), graph.name(), "m", ""); + ranges.keySet().removeIf(key -> key.startsWith(prefix)); + + PropertyKey age = schema.propertyKey("age").asInt().create(); + Assert.assertTrue(String.format("id %s reused after %s", age.id(), name.id()), + age.id().asLong() > name.id().asLong()); + + graph.clearBackend(); + } finally { + destroyGraphs(ImmutableList.of(graph)); + } + } + @Test public void testCreateMultiGraphs() { List graphs = openGraphs("g_1", NAME48); From c26d9613069419f3d8f1276ebf1c75894251c34b Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Sun, 27 Sep 2026 18:07:17 +0530 Subject: [PATCH 6/9] fix(server): stop the graphspace auth clear at the AUTH segment clearGraphSpace() deletes HUGEGRAPH/{cluster}/GRAPHSPACE/{gs}/AUTH by prefix before it drops the graphs. Graph keys sit beside AUTH under the graphspace, so the prefix also matched every key of a graph whose name starts with AUTH, such as AUTHx, and wiped its schema and config before the graph's own drop ran. End the auth prefix with the path delimiter. --- .../meta/managers/AuthMetaManager.java | 7 +++++-- .../meta/managers/AuthMetaManagerTest.java | 19 +++++++++++++++++++ 2 files changed, 24 insertions(+), 2 deletions(-) diff --git a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/AuthMetaManager.java b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/AuthMetaManager.java index bf801c0849..d8521348e6 100644 --- a/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/AuthMetaManager.java +++ b/hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/AuthMetaManager.java @@ -1045,13 +1045,16 @@ private String userListKey() { } private String authPrefix(String graphSpace) { - // HUGEGRAPH/{cluster}/GRAPHSPACE/{graphSpace}/AUTH + // HUGEGRAPH/{cluster}/GRAPHSPACE/{graphSpace}/AUTH/ + // Graph keys sit beside AUTH under the graphspace, so without the trailing + // delimiter the prefix would also match a graph named like "AUTHx" return String.join(META_PATH_DELIMITER, META_PATH_HUGEGRAPH, this.cluster, META_PATH_GRAPHSPACE, graphSpace, - META_PATH_AUTH); + META_PATH_AUTH, + ""); } private String groupKey(String group) { diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/meta/managers/AuthMetaManagerTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/meta/managers/AuthMetaManagerTest.java index 074753308e..ccc7470e23 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/meta/managers/AuthMetaManagerTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/meta/managers/AuthMetaManagerTest.java @@ -25,6 +25,7 @@ import org.apache.hugegraph.testutil.Assert; import org.apache.hugegraph.util.JsonUtil; import org.junit.Test; +import org.mockito.ArgumentCaptor; import org.mockito.Mockito; public class AuthMetaManagerTest { @@ -70,6 +71,24 @@ public void testDeleteTargetRejectsMismatchBeforeMutation() { Mockito.verify(driver, Mockito.never()).delete(Mockito.anyString()); } + @Test + public void testClearGraphAuthKeepsGraphsNamedLikeAuth() { + MetaDriver driver = Mockito.mock(MetaDriver.class); + AuthMetaManager manager = new AuthMetaManager(driver, "cluster"); + + manager.clearGraphAuth("SPACE_A"); + + ArgumentCaptor prefix = ArgumentCaptor.forClass(String.class); + Mockito.verify(driver).deleteWithPrefix(prefix.capture()); + Assert.assertEquals("HUGEGRAPH/cluster/GRAPHSPACE/SPACE_A/AUTH/", + prefix.getValue()); + // A graph named AUTHx keeps its keys directly under the graphspace + Assert.assertFalse("HUGEGRAPH/cluster/GRAPHSPACE/SPACE_A/AUTHx/SCHEMA/" + .startsWith(prefix.getValue())); + Assert.assertTrue("HUGEGRAPH/cluster/GRAPHSPACE/SPACE_A/AUTH/ROLE/r1" + .startsWith(prefix.getValue())); + } + private static HugeTarget target(String graphSpace) { HugeTarget target = new HugeTarget("target", "hugegraph", "url"); target.graphSpace(graphSpace); From 2d3dd41a8487c3a8e913f9716dca3b22d489aaa0 Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Sun, 27 Sep 2026 18:07:17 +0530 Subject: [PATCH 7/9] fix(pd): wait for a raft ReadIndex before reading an id counter IdMetaStore.getId() reads the counter from the local RocksDB and then writes the next value through raft. The read is outside raft, so a leader elected a moment ago could read the counter before it had applied its predecessor's last write and hand out the same id range twice. getId() now waits for a jraft ReadIndex first: the read index is confirmed by a quorum in the current term, and the wait ends once the local applied index reaches it. EAGAIN (no entry of the new term committed yet) and EBUSY (leadership transfer) are retried 5 times, 20 ms apart; the wait is bounded by raft.rpc-timeout. Stores that are not replicated skip the wait. --- .../apache/hugegraph/pd/meta/IdMetaStore.java | 4 + .../apache/hugegraph/pd/raft/RaftEngine.java | 73 +++++++++ .../apache/hugegraph/pd/store/HgKVStore.java | 7 + .../hugegraph/pd/store/RaftKVStore.java | 9 ++ .../hugegraph/pd/core/PDCoreSuiteTest.java | 4 + .../core/meta/IdMetaStoreReadIndexTest.java | 103 ++++++++++++ .../pd/raft/RaftEngineReadIndexTest.java | 151 ++++++++++++++++++ 7 files changed, 351 insertions(+) create mode 100644 hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/meta/IdMetaStoreReadIndexTest.java create mode 100644 hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftEngineReadIndexTest.java diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/meta/IdMetaStore.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/meta/IdMetaStore.java index 8e1fde67d6..f8e3aa391a 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/meta/IdMetaStore.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/meta/IdMetaStore.java @@ -78,6 +78,10 @@ public long getId(String key, int delta) throws PDException { Object probableLock = getLock(key); byte[] keyBs = (ID_PREFIX + key).getBytes(Charset.defaultCharset()); synchronized (probableLock) { + // The read and the put are two steps, not one raft entry: a leader elected a + // moment ago may not have applied its predecessor's last put yet, and would + // hand out the same range again without this wait + getStore().waitReadIndex(); byte[] bs = getOne(keyBs); long current = bs != null ? bytesToLong(bs) : 0L; long next = current + delta; diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftEngine.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftEngine.java index d39384909d..1258e8a88e 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftEngine.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftEngine.java @@ -46,6 +46,7 @@ import com.alipay.sofa.jraft.Node; import com.alipay.sofa.jraft.RaftGroupService; import com.alipay.sofa.jraft.Status; +import com.alipay.sofa.jraft.closure.ReadIndexClosure; import com.alipay.sofa.jraft.conf.Configuration; import com.alipay.sofa.jraft.core.Replicator; import com.alipay.sofa.jraft.core.State; @@ -74,6 +75,15 @@ public class RaftEngine { */ private static final long ALIVE_PEERS_REFRESH_MS = 1000L; + private static final int READ_INDEX_RETRIES = 5; + private static final long READ_INDEX_RETRY_DELAY_MS = 20L; + private static final ScheduledExecutorService READ_INDEX_RETRY = + Executors.newSingleThreadScheduledExecutor(runnable -> { + Thread thread = new Thread(runnable, "pd-raft-read-index-retry"); + thread.setDaemon(true); + return thread; + }); + private volatile static RaftEngine instance = new RaftEngine(); private RaftStateMachine stateMachine; private String groupId = "pd_raft"; @@ -602,6 +612,69 @@ public Node getRaftNode() { return raftNode; } + /** + * Wait until this node has applied every entry committed before the call, so a local + * read that follows sees every write acknowledged before it. jraft's ReadIndex confirms + * the commit index with a quorum in the current term and runs the closure once the + * applied index reaches it; a leader elected a moment ago first waits for an entry of + * its own term to commit. Bounded by the raft rpc timeout. + *

+ * Never call it on the state machine thread: the closure waits for that thread to apply. + */ + public void waitReadIndex() throws PDException { + waitReadIndex(this.config.getRpcTimeout()); + } + + void waitReadIndex(long timeoutMs) throws PDException { + CompletableFuture future = new CompletableFuture<>(); + readIndex(future, READ_INDEX_RETRIES); + try { + future.get(timeoutMs, TimeUnit.MILLISECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new PDException(Pdpb.ErrorType.UNKNOWN_VALUE, + "Interrupted while waiting for the raft read index", e); + } catch (TimeoutException e) { + throw new PDException(Pdpb.ErrorType.UNKNOWN_VALUE, + String.format("Raft read index timed out after %d ms", + timeoutMs)); + } catch (ExecutionException e) { + if (e.getCause() instanceof PDException) { + throw (PDException) e.getCause(); + } + throw new PDException(Pdpb.ErrorType.UNKNOWN_VALUE, e.getCause()); + } + } + + private void readIndex(CompletableFuture future, int retries) { + Node node = this.raftNode; + if (node == null) { + future.completeExceptionally(new PDException(Pdpb.ErrorType.NOT_LEADER_VALUE, + "Raft node is not started")); + return; + } + node.readIndex(new byte[0], new ReadIndexClosure() { + @Override + public void run(Status status, long index, byte[] reqCtx) { + if (status.isOk()) { + future.complete(null); + return; + } + RaftError error = status.getRaftError(); + // EAGAIN: no entry of the current term committed yet; EBUSY: transferring + if (retries > 0 && (error == RaftError.EAGAIN || error == RaftError.EBUSY)) { + READ_INDEX_RETRY.schedule(() -> readIndex(future, retries - 1), + READ_INDEX_RETRY_DELAY_MS, TimeUnit.MILLISECONDS); + return; + } + int type = error == RaftError.EPERM ? Pdpb.ErrorType.NOT_LEADER_VALUE : + Pdpb.ErrorType.UNKNOWN_VALUE; + future.completeExceptionally( + new PDException(type, "Raft read index failed: " + status)); + } + }); + } + private boolean peerEquals(PeerId p1, PeerId p2) { if (p1 == null && p2 == null) { return true; diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/store/HgKVStore.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/store/HgKVStore.java index 263cb70b28..de8e2045f5 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/store/HgKVStore.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/store/HgKVStore.java @@ -56,4 +56,11 @@ public interface HgKVStore { List scanRange(byte[] start, byte[] end); void close(); + + /** + * Wait until a local read sees every write committed before the call. A store that + * is not replicated has nothing to wait for. + */ + default void waitReadIndex() throws PDException { + } } diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/store/RaftKVStore.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/store/RaftKVStore.java index b61f07ac1d..7628b7c0f8 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/store/RaftKVStore.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/store/RaftKVStore.java @@ -179,6 +179,15 @@ public void close() { store.close(); } + /** + * Reads above are served from the local state, which on a new leader can trail the + * entries its predecessor committed; this waits for them to be applied + */ + @Override + public void waitReadIndex() throws PDException { + this.engine.waitReadIndex(); + } + /** * Need to walk the real operation of Raft */ diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java index bdacf7d371..1d108b4ce2 100644 --- a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java @@ -17,11 +17,13 @@ package org.apache.hugegraph.pd.core; +import org.apache.hugegraph.pd.core.meta.IdMetaStoreReadIndexTest; import org.apache.hugegraph.pd.core.meta.MetadataKeyHelperTest; import org.apache.hugegraph.pd.core.store.HgKVStoreImplTest; import org.apache.hugegraph.pd.raft.IpAuthHandlerTest; import org.apache.hugegraph.pd.raft.RaftEngineIpAuthIntegrationTest; import org.apache.hugegraph.pd.raft.RaftEngineLeaderAddressTest; +import org.apache.hugegraph.pd.raft.RaftEngineReadIndexTest; import org.apache.hugegraph.pd.raft.RaftEngineReadinessTest; import org.junit.runner.RunWith; import org.junit.runners.Suite; @@ -31,6 +33,7 @@ @RunWith(Suite.class) @Suite.SuiteClasses({ MetadataKeyHelperTest.class, + IdMetaStoreReadIndexTest.class, HgKVStoreImplTest.class, PDConfigTest.class, ConfigServiceTest.class, @@ -45,6 +48,7 @@ RaftEngineIpAuthIntegrationTest.class, RaftEngineLeaderAddressTest.class, RaftEngineReadinessTest.class, + RaftEngineReadIndexTest.class, // StoreNodeServiceTest.class, }) @Slf4j diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/meta/IdMetaStoreReadIndexTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/meta/IdMetaStoreReadIndexTest.java new file mode 100644 index 0000000000..5d3afa9b8e --- /dev/null +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/meta/IdMetaStoreReadIndexTest.java @@ -0,0 +1,103 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hugegraph.pd.core.meta; + +import org.apache.hugegraph.pd.common.PDException; +import org.apache.hugegraph.pd.config.PDConfig; +import org.apache.hugegraph.pd.grpc.Pdpb; +import org.apache.hugegraph.pd.meta.IdMetaStore; +import org.apache.hugegraph.pd.meta.MetadataFactory; +import org.apache.hugegraph.pd.raft.RaftEngine; +import org.apache.hugegraph.pd.store.HgKVStore; +import org.apache.hugegraph.pd.store.RaftKVStore; +import org.apache.hugegraph.testutil.Whitebox; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.mockito.InOrder; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * {@link IdMetaStore#getId} reads the counter and writes it back as two steps, so the read + * must follow a raft ReadIndex: a leader elected a moment ago may not have applied its + * predecessor's last write yet, and would hand out the same range twice. + */ +public class IdMetaStoreReadIndexTest { + + private HgKVStore originalStore; + private HgKVStore store; + private IdMetaStore idMetaStore; + + @Before + public void setUp() { + // The metadata stores share one factory-held store; swap in a mock for this test + this.originalStore = Whitebox.getInternalState(MetadataFactory.class, "store"); + this.store = mock(HgKVStore.class); + Whitebox.setInternalState(MetadataFactory.class, "store", this.store); + this.idMetaStore = new IdMetaStore(new PDConfig()); + } + + @After + public void tearDown() { + Whitebox.setInternalState(MetadataFactory.class, "store", this.originalStore); + } + + @Test + public void testGetIdReadsOnlyAfterTheReadIndex() throws PDException { + when(this.store.get(any())).thenReturn(IdMetaStore.longToBytes(100L)); + + Assert.assertEquals(100L, this.idMetaStore.getId("read-index", 10)); + + InOrder order = inOrder(this.store); + order.verify(this.store).waitReadIndex(); + order.verify(this.store).get(any()); + order.verify(this.store).put(any(), any()); + } + + @Test + public void testGetIdLeavesTheCounterAloneWhenTheReadIndexFails() throws PDException { + doThrow(new PDException(Pdpb.ErrorType.NOT_LEADER_VALUE, "not leader")) + .when(this.store).waitReadIndex(); + + PDException e = Assert.assertThrows(PDException.class, () -> { + this.idMetaStore.getId("read-index", 10); + }); + + Assert.assertEquals(Pdpb.ErrorType.NOT_LEADER_VALUE, e.getErrorCode()); + verify(this.store, never()).get(any()); + verify(this.store, never()).put(any(), any()); + } + + @Test + public void testRaftStoreWaitsOnTheRaftEngine() throws PDException { + RaftEngine engine = mock(RaftEngine.class); + HgKVStore local = mock(HgKVStore.class); + + new RaftKVStore(engine, local).waitReadIndex(); + + verify(engine).waitReadIndex(); + } +} diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftEngineReadIndexTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftEngineReadIndexTest.java new file mode 100644 index 0000000000..bb53134b4a --- /dev/null +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftEngineReadIndexTest.java @@ -0,0 +1,151 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hugegraph.pd.raft; + +import java.util.ArrayDeque; +import java.util.Arrays; +import java.util.Deque; + +import org.apache.hugegraph.pd.common.PDException; +import org.apache.hugegraph.pd.grpc.Pdpb; +import org.apache.hugegraph.testutil.Whitebox; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import com.alipay.sofa.jraft.Node; +import com.alipay.sofa.jraft.Status; +import com.alipay.sofa.jraft.closure.ReadIndexClosure; +import com.alipay.sofa.jraft.error.RaftError; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +/** + * Covers {@link RaftEngine#waitReadIndex()}, the barrier a leader runs before a local + * read-then-write such as the id counter. The raft node is a mock that answers each + * ReadIndex call the way jraft does: it sets the result and runs the closure. + */ +public class RaftEngineReadIndexTest { + + // Queued in place of a status: the closure is never run, as when no quorum answers + private static final Status NO_ANSWER = new Status(RaftError.UNKNOWN, "no answer"); + + private Node originalRaftNode; + private Node mockNode; + // One status per readIndex call, OK once the queue is empty + private final Deque answers = new ArrayDeque<>(); + + @Before + public void setUp() { + RaftEngine engine = RaftEngine.getInstance(); + this.originalRaftNode = engine.getRaftNode(); + this.mockNode = mock(Node.class); + doAnswer(invocation -> { + ReadIndexClosure closure = invocation.getArgument(1); + Status status = this.answers.isEmpty() ? Status.OK() : this.answers.poll(); + if (status != NO_ANSWER) { + closure.setResult(status.isOk() ? 42L : ReadIndexClosure.INVALID_LOG_INDEX, + invocation.getArgument(0)); + closure.run(status); + } + return null; + }).when(this.mockNode).readIndex(any(byte[].class), any(ReadIndexClosure.class)); + Whitebox.setInternalState(engine, "raftNode", this.mockNode); + } + + @After + public void tearDown() { + Whitebox.setInternalState(RaftEngine.getInstance(), "raftNode", this.originalRaftNode); + } + + @Test + public void testReturnsOnceTheReadIndexIsApplied() throws PDException { + RaftEngine.getInstance().waitReadIndex(1000L); + + verify(this.mockNode, times(1)).readIndex(any(byte[].class), + any(ReadIndexClosure.class)); + } + + @Test + public void testRetriesWhileTheNewLeaderHasNoCommitInItsTerm() throws PDException { + this.answers.addAll(Arrays.asList(new Status(RaftError.EAGAIN, "no commit yet"), + new Status(RaftError.EBUSY, "transferring"), + Status.OK())); + + RaftEngine.getInstance().waitReadIndex(1000L); + + verify(this.mockNode, times(3)).readIndex(any(byte[].class), + any(ReadIndexClosure.class)); + } + + @Test + public void testRetriesAreBounded() { + for (int i = 0; i < 10; i++) { + this.answers.add(new Status(RaftError.EAGAIN, "no commit yet")); + } + + PDException e = Assert.assertThrows(PDException.class, () -> { + RaftEngine.getInstance().waitReadIndex(5000L); + }); + + Assert.assertEquals(Pdpb.ErrorType.UNKNOWN_VALUE, e.getErrorCode()); + // The first call and five retries + verify(this.mockNode, times(6)).readIndex(any(byte[].class), + any(ReadIndexClosure.class)); + } + + @Test + public void testNodeThatLostLeadershipFailsAsNotLeader() { + this.answers.add(new Status(RaftError.EPERM, "not leader")); + + PDException e = Assert.assertThrows(PDException.class, () -> { + RaftEngine.getInstance().waitReadIndex(1000L); + }); + + Assert.assertEquals(Pdpb.ErrorType.NOT_LEADER_VALUE, e.getErrorCode()); + verify(this.mockNode, times(1)).readIndex(any(byte[].class), + any(ReadIndexClosure.class)); + } + + @Test + public void testWaitIsBoundedWhenNoQuorumAnswers() { + this.answers.add(NO_ANSWER); + + long start = System.nanoTime(); + Assert.assertThrows(PDException.class, () -> { + RaftEngine.getInstance().waitReadIndex(100L); + }); + Assert.assertTrue(System.nanoTime() - start < 1_000_000_000L); + } + + @Test + public void testMissingRaftNodeFailsAsNotLeader() { + Whitebox.setInternalState(RaftEngine.getInstance(), "raftNode", null); + + PDException e = Assert.assertThrows(PDException.class, () -> { + RaftEngine.getInstance().waitReadIndex(1000L); + }); + + Assert.assertEquals(Pdpb.ErrorType.NOT_LEADER_VALUE, e.getErrorCode()); + } +} From 1d4335fdd5813829ba8cb9464a96f5c9c9a291d6 Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Sun, 27 Sep 2026 18:07:17 +0530 Subject: [PATCH 8/9] fix(pd): take the snapshot checkpoint on the state machine thread jraft calls onSnapshotSave() on the state machine thread with the snapshot index set to the last applied index, and applies the next entry as soon as it returns. PD handed the whole save to a job thread, so the RocksDB checkpoint could include entries past the snapshot index, and a node installing that snapshot applied them a second time. Take the checkpoint before returning; only the compression stays on the job thread. The checkpoint is hard links under the store's write lock. A failed checkpoint now completes the snapshot once with an error; it used to report the error and then continue and report success as well. --- .../hugegraph/pd/raft/RaftStateMachine.java | 68 +++++--- .../hugegraph/pd/core/PDCoreSuiteTest.java | 2 + .../pd/raft/RaftStateMachineSnapshotTest.java | 163 ++++++++++++++++++ 3 files changed, 205 insertions(+), 28 deletions(-) create mode 100644 hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachineSnapshotTest.java diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java index aab90e5331..164a8cbe20 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachine.java @@ -195,45 +195,57 @@ public void onConfigurationCommitted(final Configuration conf) { log.info("Raft onConfigurationCommitted {}", conf); } + /** + * jraft calls this on the state machine thread, with the snapshot index set to the last + * applied index. The RocksDB checkpoint is taken here, before the thread applies the next + * entry: a checkpoint taken later holds entries past the snapshot index, and a node that + * installs the snapshot applies them a second time. Only the compression runs on the job + * thread. The checkpoint is hard links under the store's write lock, so it is short. + */ @Override public void onSnapshotSave(final SnapshotWriter writer, final Closure done) { - MetadataService.getUninterruptibleJobs().submit(() -> { - lock.lock(); + String snapshotDir = writer.getPath() + File.separator + SNAPSHOT_DIR_NAME; + lock.lock(); + try { + log.info("start snapshot save"); try { - log.info("start snapshot save"); - String snapshotDir = writer.getPath() + File.separator + SNAPSHOT_DIR_NAME; + FileUtils.deleteDirectory(new File(snapshotDir)); + FileUtils.forceMkdir(new File(snapshotDir)); + } catch (IOException e) { + log.error("Failed to create snapshot directory {}", snapshotDir); + done.run(new Status(RaftError.EIO, e.toString())); + return; + } + for (RaftTaskHandler taskHandler : taskHandlers) { try { - FileUtils.deleteDirectory(new File(snapshotDir)); - FileUtils.forceMkdir(new File(snapshotDir)); - } catch (IOException e) { - log.error("Failed to create snapshot directory {}", snapshotDir); + KVOperation op = KVOperation.createSaveSnapshot(snapshotDir); + taskHandler.invoke(op, null); + log.info("Raft onSnapshotSave success"); + } catch (PDException e) { + log.error("Raft onSnapshotSave failed. {}", e.toString()); done.run(new Status(RaftError.EIO, e.toString())); return; } - for (RaftTaskHandler taskHandler : taskHandlers) { - try { - KVOperation op = KVOperation.createSaveSnapshot(snapshotDir); - taskHandler.invoke(op, null); - log.info("Raft onSnapshotSave success"); - } catch (PDException e) { - log.error("Raft onSnapshotSave failed. {}", e.toString()); - done.run(new Status(RaftError.EIO, e.toString())); - } - } + } + } catch (Exception e) { + log.error("failed to save snapshot", e); + done.run(new Status(RaftError.EIO, e.toString())); + return; + } finally { + lock.unlock(); + } + + MetadataService.getUninterruptibleJobs().submit(() -> { + lock.lock(); + try { // compress - try { - compressSnapshot(writer); - FileUtils.deleteDirectory(new File(snapshotDir)); - } catch (Exception e) { - log.error("Failed to delete snapshot directory {}, {}", snapshotDir, - e.toString()); - done.run(new Status(RaftError.EIO, e.toString())); - return; - } + compressSnapshot(writer); + FileUtils.deleteDirectory(new File(snapshotDir)); done.run(Status.OK()); log.info("snapshot save done"); } catch (Exception e) { - log.error("failed to save snapshot", e); + log.error("Failed to compress snapshot directory {}, {}", snapshotDir, + e.toString()); done.run(new Status(RaftError.EIO, e.toString())); } finally { lock.unlock(); diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java index 1d108b4ce2..ccad7ef0fd 100644 --- a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java @@ -25,6 +25,7 @@ import org.apache.hugegraph.pd.raft.RaftEngineLeaderAddressTest; import org.apache.hugegraph.pd.raft.RaftEngineReadIndexTest; import org.apache.hugegraph.pd.raft.RaftEngineReadinessTest; +import org.apache.hugegraph.pd.raft.RaftStateMachineSnapshotTest; import org.junit.runner.RunWith; import org.junit.runners.Suite; @@ -49,6 +50,7 @@ RaftEngineLeaderAddressTest.class, RaftEngineReadinessTest.class, RaftEngineReadIndexTest.class, + RaftStateMachineSnapshotTest.class, // StoreNodeServiceTest.class, }) @Slf4j diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachineSnapshotTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachineSnapshotTest.java new file mode 100644 index 0000000000..aa944f7c83 --- /dev/null +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftStateMachineSnapshotTest.java @@ -0,0 +1,163 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hugegraph.pd.raft; + +import java.io.File; +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.commons.io.FileUtils; +import org.apache.hugegraph.pd.common.PDException; +import org.apache.hugegraph.pd.grpc.Pdpb; +import org.apache.hugegraph.pd.service.MetadataService; +import org.apache.hugegraph.testutil.Whitebox; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import com.alipay.sofa.jraft.Closure; +import com.alipay.sofa.jraft.Status; +import com.alipay.sofa.jraft.error.RaftError; +import com.alipay.sofa.jraft.storage.snapshot.SnapshotWriter; +import com.google.protobuf.Message; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * jraft calls {@code onSnapshotSave} on the state machine thread with the snapshot index + * set to the last applied index, and applies the next entry once it returns. The checkpoint + * must therefore be taken before the call returns: a checkpoint taken later on a job thread + * holds entries past the snapshot index, which a node installing it applies a second time. + */ +public class RaftStateMachineSnapshotTest { + + private ThreadPoolExecutor originalJobs; + private ThreadPoolExecutor jobs; + private File snapshotPath; + private SnapshotWriter writer; + + @Before + public void setUp() throws IOException { + this.originalJobs = MetadataService.getUninterruptibleJobs(); + this.jobs = new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS, + new LinkedBlockingQueue<>()); + Whitebox.setInternalState(MetadataService.class, "uninterruptibleJobs", this.jobs); + + this.snapshotPath = Files.createTempDirectory("pd-snapshot-save").toFile(); + this.writer = mock(SnapshotWriter.class); + when(this.writer.getPath()).thenReturn(this.snapshotPath.getAbsolutePath()); + when(this.writer.addFile(anyString(), any(Message.class))).thenReturn(true); + } + + @After + public void tearDown() throws IOException { + Whitebox.setInternalState(MetadataService.class, "uninterruptibleJobs", + this.originalJobs); + this.jobs.shutdownNow(); + FileUtils.deleteDirectory(this.snapshotPath); + } + + @Test + public void testCheckpointIsTakenBeforeTheNextEntryIsApplied() throws Exception { + AtomicLong applied = new AtomicLong(10L); + AtomicReference checkpointThread = new AtomicReference<>(); + AtomicLong checkpointed = new AtomicLong(-1L); + RaftStateMachine machine = new RaftStateMachine(); + machine.addTaskHandler((op, response) -> { + if (op.getOp() == KVOperation.SAVE_SNAPSHOT) { + checkpointThread.set(Thread.currentThread()); + checkpointed.set(applied.get()); + writeCheckpoint((String) op.getAttach(), applied.get()); + } + return false; + }); + RecordingClosure done = new RecordingClosure(); + + machine.onSnapshotSave(this.writer, done); + // What the state machine thread does next: apply the entry after the snapshot + applied.incrementAndGet(); + + Assert.assertSame(Thread.currentThread(), checkpointThread.get()); + Assert.assertEquals(10L, checkpointed.get()); + Assert.assertTrue(done.await()); + Assert.assertEquals(1, done.statuses.size()); + Assert.assertTrue(done.statuses.get(0).isOk()); + Assert.assertTrue(new File(this.snapshotPath, "snapshot.zip").isFile()); + } + + @Test + public void testFailedCheckpointCompletesTheSnapshotOnceWithAnError() throws Exception { + RaftStateMachine machine = new RaftStateMachine(); + machine.addTaskHandler((op, response) -> { + if (op.getOp() == KVOperation.SAVE_SNAPSHOT) { + throw new PDException(Pdpb.ErrorType.ROCKSDB_SAVE_SNAPSHOT_ERROR_VALUE, + "checkpoint failed"); + } + return false; + }); + RecordingClosure done = new RecordingClosure(); + + machine.onSnapshotSave(this.writer, done); + + Assert.assertTrue(done.await()); + // Drain the job thread so a second completion, if any, has run + this.jobs.submit(() -> { }).get(10, TimeUnit.SECONDS); + Assert.assertEquals(1, done.statuses.size()); + Assert.assertEquals(RaftError.EIO, done.statuses.get(0).getRaftError()); + Assert.assertFalse(new File(this.snapshotPath, "snapshot.zip").exists()); + } + + private static void writeCheckpoint(String dir, long index) { + try { + FileUtils.forceMkdir(new File(dir)); + Files.write(new File(dir, "CURRENT").toPath(), + Long.toString(index).getBytes(StandardCharsets.UTF_8)); + } catch (IOException e) { + throw new IllegalStateException(e); + } + } + + private static final class RecordingClosure implements Closure { + + private final List statuses = new CopyOnWriteArrayList<>(); + private final CountDownLatch latch = new CountDownLatch(1); + + @Override + public void run(Status status) { + this.statuses.add(status); + this.latch.countDown(); + } + + private boolean await() throws InterruptedException { + return this.latch.await(10, TimeUnit.SECONDS); + } + } +} From 8d6e18684ef37cbf96e948974101b5ca7c8e95d4 Mon Sep 17 00:00:00 2001 From: Himanshu Verma Date: Sun, 27 Sep 2026 18:07:17 +0530 Subject: [PATCH 9/9] fix(pd): give the raft connect timeout its own option PD set jraft's connect, default and install-snapshot rpc timeouts all from raft.rpc-timeout (10 s). A candidate opens a connection to each peer while it holds the node lock, and a peer that accepts the TCP connection but never answers the ping (a stopped process, a lost host) stalls the election for the whole connect timeout. Losing the host of the first leader after cluster start left the PD cluster without a leader for about 11 s. Add raft.rpc-connect-timeout (default 1000 ms, jraft's own default). Split out raft.rpc-install-snapshot-timeout as well, so the connect timeout can be lowered without shortening snapshot installs; its default 0 keeps the current value, raft.rpc-timeout. raft.rpc-timeout keeps its meaning for other raft requests. --- hugegraph-pd/docs/configuration.md | 3 + .../apache/hugegraph/pd/config/PDConfig.java | 13 +++ .../apache/hugegraph/pd/raft/RaftEngine.java | 18 ++++- .../static/conf/application.yml.template | 6 ++ .../src/main/resources/application.yml | 6 ++ .../hugegraph/pd/core/PDCoreSuiteTest.java | 2 + .../pd/raft/RaftEngineRpcTimeoutTest.java | 79 +++++++++++++++++++ 7 files changed, 124 insertions(+), 3 deletions(-) create mode 100644 hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftEngineRpcTimeoutTest.java diff --git a/hugegraph-pd/docs/configuration.md b/hugegraph-pd/docs/configuration.md index 4e60d8e8be..059e15126e 100644 --- a/hugegraph-pd/docs/configuration.md +++ b/hugegraph-pd/docs/configuration.md @@ -119,6 +119,9 @@ raft: |-----------|------|---------|-------------| | `raft.address` | String | `127.0.0.1:8610` | Raft service address for this PD node. Format: `:`. Must be unique across all PD nodes. | | `raft.peers-list` | String | `127.0.0.1:8610` | Comma-separated list of all PD nodes' Raft addresses. Used for cluster formation and leader election. | +| `raft.rpc-timeout` | Integer | `10000` | Timeout of a Raft RPC between PD nodes, in milliseconds. | +| `raft.rpc-connect-timeout` | Integer | `1000` | Timeout for opening a connection to a peer, in milliseconds. A candidate opens its connections one peer at a time while it holds the Raft node lock, so an election waits this long for each peer that accepts the connection but does not answer (a stopped process, a lost host). | +| `raft.rpc-install-snapshot-timeout` | Integer | `0` | Timeout for sending a snapshot to a follower that fell behind the log, in milliseconds. It bounds the transfer of the whole metadata store. `0` uses `raft.rpc-timeout`, the value this timeout had before it became a separate option; raise it if large snapshots time out. | **Critical Rules**: 1. `raft.address` must be unique for each PD node diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/config/PDConfig.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/config/PDConfig.java index 56bd58b34d..3049915602 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/config/PDConfig.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/config/PDConfig.java @@ -186,11 +186,24 @@ public class Raft { private int snapshotInterval; @Value("${raft.rpc-timeout:10000}") private int rpcTimeout; + // Bounds the ping that opens a connection to a peer; kept apart from rpc-timeout + // because a candidate opens its connections while it holds the raft node lock + @Value("${raft.rpc-connect-timeout:1000}") + private int rpcConnectTimeout = 1000; + // A follower that fell behind the log receives the whole store within this time; + // 0 keeps the timeout it always had, raft.rpc-timeout + @Value("${raft.rpc-install-snapshot-timeout:0}") + private int rpcInstallSnapshotTimeout; @Value("${grpc.host}") private String host; @Value("${server.port}") private int port; + public int getRpcInstallSnapshotTimeout() { + return this.rpcInstallSnapshotTimeout > 0 ? this.rpcInstallSnapshotTimeout : + this.rpcTimeout; + } + @Value("${pd.cluster_id:1}") private long clusterId; @Value("${grpc.port}") diff --git a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftEngine.java b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftEngine.java index 1258e8a88e..59a90c6686 100644 --- a/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftEngine.java +++ b/hugegraph-pd/hg-pd-core/src/main/java/org/apache/hugegraph/pd/raft/RaftEngine.java @@ -143,9 +143,7 @@ public synchronized boolean init(PDConfig.Raft config) { // Snapshot interval nodeOptions.setSnapshotIntervalSecs(config.getSnapshotInterval()); - nodeOptions.setRpcConnectTimeoutMs(config.getRpcTimeout()); - nodeOptions.setRpcDefaultTimeout(config.getRpcTimeout()); - nodeOptions.setRpcInstallSnapshotTimeout(config.getRpcTimeout()); + setRpcTimeouts(nodeOptions, config); // TODO: tune RaftOptions for PD (see hugegraph-store PartitionEngine for reference) final PeerId serverId = JRaftUtils.getPeerId(config.getAddress()); @@ -162,6 +160,20 @@ public synchronized boolean init(PDConfig.Raft config) { return this.raftNode != null; } + /** + * Wire the three raft rpc timeouts from their own options. jraft pings a peer to open a + * connection, and a candidate does so for every peer while it holds the node lock, so a + * peer that accepts the connection but never answers (a stopped process, a lost host) + * stalls the election, and the answers to the other peers' votes, for the whole connect + * timeout. Installing a snapshot sends the whole store and needs far longer than a + * normal request. + */ + static void setRpcTimeouts(NodeOptions nodeOptions, PDConfig.Raft config) { + nodeOptions.setRpcConnectTimeoutMs(config.getRpcConnectTimeout()); + nodeOptions.setRpcDefaultTimeout(config.getRpcTimeout()); + nodeOptions.setRpcInstallSnapshotTimeout(config.getRpcInstallSnapshotTimeout()); + } + /** * Create a Raft RPC Server for communication between PDs */ diff --git a/hugegraph-pd/hg-pd-dist/src/assembly/static/conf/application.yml.template b/hugegraph-pd/hg-pd-dist/src/assembly/static/conf/application.yml.template index 43bdf99fc4..1fc2fc4fea 100644 --- a/hugegraph-pd/hg-pd-dist/src/assembly/static/conf/application.yml.template +++ b/hugegraph-pd/hg-pd-dist/src/assembly/static/conf/application.yml.template @@ -63,6 +63,12 @@ raft: address: $RAFT_ADDRESS$ # raft cluster peers-list: $RAFT_PEERS_LIST$ + # Timeout for opening a connection to a peer, in milliseconds; an election waits this + # long for each peer that accepts the connection but does not answer + rpc-connect-timeout: 1000 + # Timeout for sending a snapshot to a follower that fell behind the log, in milliseconds; + # 0 uses rpc-timeout + rpc-install-snapshot-timeout: 0 # The interval between snapshot generation, in seconds snapshotInterval: 300 metrics: true diff --git a/hugegraph-pd/hg-pd-service/src/main/resources/application.yml b/hugegraph-pd/hg-pd-service/src/main/resources/application.yml index f8c1ecea84..fb44a3a4c1 100644 --- a/hugegraph-pd/hg-pd-service/src/main/resources/application.yml +++ b/hugegraph-pd/hg-pd-service/src/main/resources/application.yml @@ -72,6 +72,12 @@ raft: peers-list: 127.0.0.1:8610 # The read and write timeout period of the raft rpc, in milliseconds rpc-timeout: 10000 + # Timeout for opening a connection to a peer, in milliseconds; an election waits this + # long for each peer that accepts the connection but does not answer + rpc-connect-timeout: 1000 + # Timeout for sending a snapshot to a follower that fell behind the log, in milliseconds; + # 0 uses rpc-timeout + rpc-install-snapshot-timeout: 0 # The interval between snapshot generation, in seconds snapshotInterval: 300 metrics: true diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java index ccad7ef0fd..a148dc5825 100644 --- a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/core/PDCoreSuiteTest.java @@ -25,6 +25,7 @@ import org.apache.hugegraph.pd.raft.RaftEngineLeaderAddressTest; import org.apache.hugegraph.pd.raft.RaftEngineReadIndexTest; import org.apache.hugegraph.pd.raft.RaftEngineReadinessTest; +import org.apache.hugegraph.pd.raft.RaftEngineRpcTimeoutTest; import org.apache.hugegraph.pd.raft.RaftStateMachineSnapshotTest; import org.junit.runner.RunWith; import org.junit.runners.Suite; @@ -51,6 +52,7 @@ RaftEngineReadinessTest.class, RaftEngineReadIndexTest.class, RaftStateMachineSnapshotTest.class, + RaftEngineRpcTimeoutTest.class, // StoreNodeServiceTest.class, }) @Slf4j diff --git a/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftEngineRpcTimeoutTest.java b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftEngineRpcTimeoutTest.java new file mode 100644 index 0000000000..3c22da393b --- /dev/null +++ b/hugegraph-pd/hg-pd-test/src/main/java/org/apache/hugegraph/pd/raft/RaftEngineRpcTimeoutTest.java @@ -0,0 +1,79 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hugegraph.pd.raft; + +import org.apache.hugegraph.pd.config.PDConfig; +import org.junit.Assert; +import org.junit.Test; +import org.springframework.beans.factory.annotation.Value; + +import com.alipay.sofa.jraft.option.NodeOptions; + +/** + * A candidate opens a connection to every peer while it holds the raft node lock, so the + * connect timeout bounds how long one unanswering peer stalls an election. It used to be + * {@code raft.rpc-timeout} (10 s), shared with the install-snapshot timeout; each now has + * its own option, and the install-snapshot timeout still follows rpc-timeout unless set. + */ +public class RaftEngineRpcTimeoutTest { + + @Test + public void testEachRpcTimeoutComesFromItsOwnOption() { + PDConfig.Raft raft = new PDConfig().new Raft(); + raft.setRpcTimeout(7000); + raft.setRpcConnectTimeout(900); + raft.setRpcInstallSnapshotTimeout(120000); + NodeOptions options = new NodeOptions(); + + RaftEngine.setRpcTimeouts(options, raft); + + Assert.assertEquals(900, options.getRpcConnectTimeoutMs()); + Assert.assertEquals(7000, options.getRpcDefaultTimeout()); + Assert.assertEquals(120000, options.getRpcInstallSnapshotTimeout()); + } + + @Test + public void testDefaults() throws NoSuchFieldException { + // The connect timeout takes jraft's default of 1 s + Assert.assertEquals("${raft.rpc-connect-timeout:1000}", + valueOf("rpcConnectTimeout")); + Assert.assertEquals("${raft.rpc-install-snapshot-timeout:0}", + valueOf("rpcInstallSnapshotTimeout")); + Assert.assertEquals("${raft.rpc-timeout:10000}", valueOf("rpcTimeout")); + + PDConfig.Raft raft = new PDConfig().new Raft(); + Assert.assertEquals(1000, raft.getRpcConnectTimeout()); + Assert.assertEquals(new NodeOptions().getRpcConnectTimeoutMs(), + raft.getRpcConnectTimeout()); + } + + @Test + public void testInstallSnapshotTimeoutFollowsRpcTimeoutWhenUnset() { + PDConfig.Raft raft = new PDConfig().new Raft(); + raft.setRpcTimeout(10000); + NodeOptions options = new NodeOptions(); + + RaftEngine.setRpcTimeouts(options, raft); + + Assert.assertEquals(10000, options.getRpcInstallSnapshotTimeout()); + } + + private static String valueOf(String field) throws NoSuchFieldException { + return PDConfig.Raft.class.getDeclaredField(field).getAnnotation(Value.class).value(); + } +}