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/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..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 @@ -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"; @@ -133,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()); @@ -152,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 */ @@ -602,6 +624,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/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-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-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 bdacf7d371..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 @@ -17,12 +17,16 @@ 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.apache.hugegraph.pd.raft.RaftEngineRpcTimeoutTest; +import org.apache.hugegraph.pd.raft.RaftStateMachineSnapshotTest; import org.junit.runner.RunWith; import org.junit.runners.Suite; @@ -31,6 +35,7 @@ @RunWith(Suite.class) @Suite.SuiteClasses({ MetadataKeyHelperTest.class, + IdMetaStoreReadIndexTest.class, HgKVStoreImplTest.class, PDConfigTest.class, ConfigServiceTest.class, @@ -45,6 +50,9 @@ RaftEngineIpAuthIntegrationTest.class, RaftEngineLeaderAddressTest.class, 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/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()); + } +} 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(); + } +} 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); + } + } +} 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..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 @@ -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(); } } @@ -315,37 +397,57 @@ 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); - // Status transitions are internal bookkeeping; notifying here causes a - // broadcast storm for every updateSchemaStatus() call from background jobs. + 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. + this.bumpSchemaVersion(); } @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 // of TASK_SYNC_DELETION (which only gates removal flows). this.notifySchemaCacheClear(); } - private void updateCache(SchemaElement schema) { - this.resetCachedAllIfReachedCapacity(); - - // update id cache + private void updateCache(SchemaElement schema, long generation) { 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); + synchronized (this.arrayCaches) { + 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 + this.idCache.update(prefixedId, schema); + + // update name cache + this.nameCache.update(prefixedName, schema); - // update optimized array cache - this.arrayCaches.updateIfNeeded(schema); + // update optimized array cache + this.arrayCaches.updateIfNeeded(schema); + } } @Override @@ -354,9 +456,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 +493,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 +562,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 +604,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 +620,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 +633,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 +763,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/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-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..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,6 +4,15 @@ gremlin.graph=org.apache.hugegraph.HugeFactory # cache config #schema.cache_capacity=100000 +# 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 vertex.cache_type=l2 #vertex.cache_capacity=10000000 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); 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); 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..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 @@ -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,317 @@ 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 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"); + 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 + 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 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, + 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; + } +}