diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/AuthenticationFilter.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/AuthenticationFilter.java index 96eb273be7..00b836e31d 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/AuthenticationFilter.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/AuthenticationFilter.java @@ -77,6 +77,7 @@ public class AuthenticationFilter implements ContainerRequestFilter, ContainerRe private static final AntPathMatcher MATCHER = new AntPathMatcher(); private static final Set FIXED_WHITE_API_SET = ImmutableSet.of( "versions", + "readiness", "openapi.json" ); /** Remove auth/login API from whitelist */ diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/LoadDetectFilter.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/LoadDetectFilter.java index 1df19f5e5c..390c2c9e80 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/LoadDetectFilter.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/LoadDetectFilter.java @@ -53,7 +53,8 @@ public class LoadDetectFilter implements ContainerRequestFilter { "", "apis", "metrics", - "versions" + "versions", + "readiness" ); // Call gc every 30+ seconds if memory is low and request frequently diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/PathFilter.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/PathFilter.java index 5e4dd5081c..016740d181 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/PathFilter.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/filter/PathFilter.java @@ -57,6 +57,7 @@ public class PathFilter implements ContainerRequestFilter { "apis", "metrics", "versions", + "readiness", "health", "gremlin", "graphs/auth", diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/profile/ReadinessAPI.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/profile/ReadinessAPI.java new file mode 100644 index 0000000000..44186db48e --- /dev/null +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/profile/ReadinessAPI.java @@ -0,0 +1,69 @@ +/* + * 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.api.profile; + +import java.util.Map; + +import org.apache.hugegraph.api.API; +import org.apache.hugegraph.config.HugeConfig; +import org.apache.hugegraph.config.ServerOptions; +import org.apache.hugegraph.core.GraphManager; +import org.apache.hugegraph.util.JsonUtil; + +import com.codahale.metrics.annotation.Timed; + +import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.annotation.security.PermitAll; +import jakarta.inject.Singleton; +import jakarta.ws.rs.GET; +import jakarta.ws.rs.Path; +import jakarta.ws.rs.Produces; +import jakarta.ws.rs.core.Context; +import jakarta.ws.rs.core.Response; + +/** + * Storage-aware readiness for Kubernetes and load balancers: on hstore 200 + * while at least one known Store answers this server, 503 while none does + * (or, before any Store list is known, while PD does not answer); on hbase + * 200 while the cluster answers an admin call within the budget. Unauthenticated, like + * /versions, so that an httpGet probe needs no credential; the body carries + * no addresses and no raw exception text. + */ +@Path("readiness") +@Singleton +@Tag(name = "ReadinessAPI") +public class ReadinessAPI extends API { + + @GET + @Timed + @Produces(APPLICATION_JSON_WITH_CHARSET) + @PermitAll + public Response get(@Context GraphManager manager, @Context HugeConfig conf) { + Map body = StorageReadiness.check( + manager, conf.get(ServerOptions.READINESS_TIMEOUT), + conf.get(ServerOptions.READINESS_CACHE_TTL), + conf.get(ServerOptions.READINESS_MAX_WAITERS)); + Response.Status status = StorageReadiness.isReady(body) ? + Response.Status.OK : + Response.Status.SERVICE_UNAVAILABLE; + return Response.status(status) + .type(APPLICATION_JSON_WITH_CHARSET) + .entity(JsonUtil.toJson(body)) + .build(); + } +} diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/profile/StorageReadiness.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/profile/StorageReadiness.java new file mode 100644 index 0000000000..47ad8c6355 --- /dev/null +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/api/profile/StorageReadiness.java @@ -0,0 +1,366 @@ +/* + * 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.api.profile; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import java.util.stream.Collectors; + +import org.apache.hugegraph.HugeGraph; +import org.apache.hugegraph.auth.HugeGraphAuthProxy; +import org.apache.hugegraph.core.GraphManager; +import org.apache.hugegraph.util.Log; +import org.slf4j.Logger; + +import com.google.common.collect.ImmutableSet; + +/** + * Whether this server can serve graph traffic, for a readiness probe. + * Graphs on an embedded backend are ready as soon as the REST layer answers. + * Graphs on a remote backend are probed through the backend's + * "storage_readiness" metadata: on hstore at least one known Store answers a + * cheap direct call (PD is only needed until the first Store list is known), + * on hbase the cluster answers an admin call about one of the graph's tables. The storage is shared by every + * hstore graph of the process, so one graph is probed and the result is + * reused for a short TTL to keep repeated probes cheap. The body never + * carries raw exception text, since the endpoint is unauthenticated. + */ +public final class StorageReadiness { + + public static final String STORAGE_READINESS_META = "storage_readiness"; + public static final String BACKEND_HSTORE = "hstore"; + public static final String BACKEND_HBASE = "hbase"; + /** Backends on a remote cluster, whose availability the probe checks. */ + public static final Set REMOTE_BACKENDS = ImmutableSet.of(BACKEND_HSTORE, + BACKEND_HBASE); + + private static final Logger LOG = Log.logger(StorageReadiness.class); + + /** A probe result with the time it was taken: the TTL check and the body come from one probe. */ + private static final class Cached { + + private final Map result; + private final long at; + + private Cached(Map result, long at) { + this.result = result; + this.at = at; + } + } + + private static final AtomicReference LAST = new AtomicReference<>(); + private static final AtomicReference>> IN_FLIGHT = + new AtomicReference<>(); + private static final AtomicInteger WAITERS = new AtomicInteger(); + /** Probes of independent backend configurations run side by side here. */ + private static final ExecutorService PROBES = + Executors.newCachedThreadPool(r -> { + Thread t = new Thread(r, "storage-readiness-probe"); + t.setDaemon(true); + return t; + }); + + private StorageReadiness() { + } + + /** One storage probe with a time budget in ms. */ + public interface Probe { + + Map probe(long timeoutMs) throws Exception; + } + + public static Map check(GraphManager manager, long timeoutMs, + long cacheTtlMs, int maxWaiters) { + // The graphs are auth proxies and the probe request carries no user, + // so look the graphs up and probe them as the internal admin, the way + // other internal paths do; the result carries no data or addresses + List> holder = new ArrayList<>(1); + HugeGraphAuthProxy.runAsAdmin(() -> { + List remotes = remoteGraphs(manager); + if (remotes.isEmpty()) { + Map body = new LinkedHashMap<>(); + body.put("ready", true); + body.put("storage", "embedded"); + body.put("reason", "no graph on a remote storage"); + holder.add(body); + return; + } + String storage = remotes.stream().map(r -> r.backend).distinct() + .collect(Collectors.joining(",")); + holder.add(check(storage, t -> probeAll(remotes, t), timeoutMs, cacheTtlMs, maxWaiters)); + }); + return holder.get(0); + } + + /** One graph per distinct remote backend configuration, in graph order; public for the unit test. */ + public static final class RemoteGraph { + + final String name; + final String backend; + final Probe probe; + + public RemoteGraph(String name, String backend, Probe probe) { + this.name = name; + this.backend = backend; + this.probe = probe; + } + } + + /** + * Readiness covers every remote backend configuration this server + * uses: graphs that share a configuration (the same PD peers, the same + * HBase hosts and namespace) share one probe; independent clusters are + * probed side by side within the common budget, and ready means all of + * them are ready. + */ + static List remoteGraphs(GraphManager manager) { + List out = new ArrayList<>(); + Set seen = new HashSet<>(); + for (String name : manager.graphs()) { + try { + HugeGraph graph = manager.graph(name); + if (graph == null || !REMOTE_BACKENDS.contains(graph.backend())) { + continue; + } + String key = graph.backend() + "|" + configKey(graph); + if (!seen.add(key)) { + continue; + } + out.add(new RemoteGraph(name, graph.backend(), probeOf(graph))); + } catch (Throwable e) { + LOG.debug("Skip graph {} while looking for remote storages", name, e); + } + } + return out; + } + + /** + * The probe of one graph. `metadata()` auto-opens the thread-local graph + * transaction (it goes through graphTransaction()), and a transaction + * left open on a request or probe thread makes the graph's close() fail + * its all-threads-closed check, so the probe closes it on the way out. + */ + public static Probe probeOf(HugeGraph graph) { + return t -> { + try { + return graph.metadata(null, STORAGE_READINESS_META, t); + } finally { + if (graph.tx().isOpen()) { + graph.tx().close(); + } + } + }; + } + + private static String configKey(HugeGraph graph) { + try { + if (BACKEND_HSTORE.equals(graph.backend())) { + return String.valueOf(graph.configuration().getString("pd.peers")); + } + return graph.configuration().getString("hbase.hosts") + "/" + + graph.configuration().getString("hbase.namespace"); + } catch (Throwable e) { + return graph.name(); + } + } + + /** Every configuration's probe in parallel on the shared budget; the body merges them. */ + public static Map probeAll(List remotes, long timeoutMs) + throws Exception { + if (remotes.size() == 1) { + Map body = new LinkedHashMap<>(remotes.get(0).probe.probe(timeoutMs)); + body.put("probes", List.of(probeEntry(remotes.get(0), body))); + return body; + } + long deadline = System.currentTimeMillis() + timeoutMs; + List>> futures = new ArrayList<>(); + for (RemoteGraph r : remotes) { + futures.add(CompletableFuture.supplyAsync(() -> { + // the auth context is a thread local of the request thread: the + // probe runs as the internal admin on its own thread as well + List> holder = new ArrayList<>(1); + HugeGraphAuthProxy.runAsAdmin(() -> { + try { + holder.add(r.probe.probe(timeoutMs)); + } catch (Exception e) { + Map failed = new LinkedHashMap<>(); + failed.put("ready", false); + failed.put("reason", "probe failed: " + e.getClass().getSimpleName()); + holder.add(failed); + } + }); + return holder.get(0); + }, PROBES)); + } + List> entries = new ArrayList<>(); + Map first = null; + boolean ready = true; + String reason = "ok"; + for (int i = 0; i < remotes.size(); i++) { + Map result; + try { + long remaining = Math.max(1L, deadline - System.currentTimeMillis()); + result = futures.get(i).get(remaining, TimeUnit.MILLISECONDS); + } catch (TimeoutException e) { + futures.get(i).cancel(true); + result = new LinkedHashMap<>(); + result.put("ready", false); + result.put("reason", "did not answer within " + timeoutMs + " ms"); + } + if (first == null) { + first = result; + } + Map entry = probeEntry(remotes.get(i), result); + entries.add(entry); + if (!isReady(result)) { + if (ready) { + // no graph name: the endpoint is unauthenticated + reason = remotes.get(i).backend + " configuration " + (i + 1) + " of " + + remotes.size() + ": " + result.get("reason"); + } + ready = false; + } + } + Map body = new LinkedHashMap<>(first); + body.put("ready", ready); + body.put("reason", reason); + body.put("probes", entries); + return body; + } + + /** One probe in the public body: backend and outcome, no graph name (unauthenticated endpoint). */ + private static Map probeEntry(RemoteGraph r, Map result) { + Map entry = new LinkedHashMap<>(); + entry.put("storage", r.backend); + entry.put("ready", isReady(result)); + entry.put("reason", result.get("reason")); + return entry; + } + + public static Map check(Probe probe, long timeoutMs, long cacheTtlMs) { + return check(BACKEND_HSTORE, probe, timeoutMs, cacheTtlMs, Integer.MAX_VALUE); + } + + /** + * At most one probe runs at a time and every caller gets its answer: + * the first caller runs the probe on its own thread, concurrent callers + * wait for that result, each bounded by its own timeout, and at most + * {@code maxWaiters} of them wait at once: the rest get an immediate + * not-ready, so a burst of probes during slow storage cannot hold the + * REST worker pool (the endpoint is unauthenticated and outside the + * load-shedding filter). No monitor is held during the probe. + */ + public static Map check(String storage, Probe probe, long timeoutMs, + long cacheTtlMs, int maxWaiters) { + Cached last = LAST.get(); + if (last != null && System.currentTimeMillis() - last.at < cacheTtlMs) { + Map body = new LinkedHashMap<>(last.result); + body.put("cached", true); + return body; + } + CompletableFuture> mine = new CompletableFuture<>(); + CompletableFuture> running = IN_FLIGHT.get(); + if (running == null && IN_FLIGHT.compareAndSet(null, mine)) { + Map body; + try { + body = probeOnce(storage, probe, timeoutMs); + LAST.set(new Cached(body, System.currentTimeMillis())); + } finally { + IN_FLIGHT.set(null); + } + mine.complete(body); + return new LinkedHashMap<>(body); + } + if (running == null) { + running = IN_FLIGHT.get(); + } + if (running == null) { + // The owner finished between our two reads: its result is cached + return check(storage, probe, timeoutMs, cacheTtlMs, maxWaiters); + } + if (WAITERS.incrementAndGet() > maxWaiters) { + WAITERS.decrementAndGet(); + return notReady(storage, "too many readiness callers waiting for the probe (" + + maxWaiters + ")"); + } + try { + Map body = new LinkedHashMap<>( + running.get(timeoutMs, TimeUnit.MILLISECONDS)); + body.put("shared", true); + return body; + } catch (TimeoutException e) { + return notReady(storage, "a probe is still running after " + timeoutMs + " ms"); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return notReady(storage, "interrupted"); + } catch (ExecutionException e) { + return notReady(storage, "probe failed: " + + e.getCause().getClass().getSimpleName()); + } finally { + WAITERS.decrementAndGet(); + } + } + + private static Map probeOnce(String storage, Probe probe, long timeoutMs) { + Map body = new LinkedHashMap<>(); + body.put("ready", false); + body.put("storage", storage); + try { + Map result = probe.probe(timeoutMs); + body.putAll(result); + } catch (Throwable e) { + LOG.warn("Storage readiness probe failed", e); + body.put("ready", false); + body.put("reason", "probe failed: " + e.getClass().getSimpleName()); + } + body.put("cached", false); + return body; + } + + private static Map notReady(String storage, String reason) { + Map body = new LinkedHashMap<>(); + body.put("ready", false); + body.put("storage", storage); + body.put("reason", reason); + body.put("cached", false); + return body; + } + + public static boolean isReady(Map body) { + return Boolean.TRUE.equals(body.get("ready")); + } + + public static void resetCache() { + LAST.set(null); + IN_FLIGHT.set(null); + WAITERS.set(0); + } + +} diff --git a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/config/ServerOptions.java b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/config/ServerOptions.java index e6ed6954a1..117427cb67 100644 --- a/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/config/ServerOptions.java +++ b/hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/config/ServerOptions.java @@ -202,6 +202,41 @@ public class ServerOptions extends OptionHolder { 300 ); + public static final ConfigOption READINESS_TIMEOUT = + new ConfigOption<>( + "readiness.timeout", + "The whole time budget in ms of one GET /readiness probe: " + + "one cheap call to every known Store in parallel, the " + + "first answer wins; the PD call that refreshes the Store " + + "list runs in the background and is only waited for " + + "before the first list is known. No Store answering " + + "within the budget makes the server report not ready.", + rangeInt(100, 60000), + 1000 + ); + + public static final ConfigOption READINESS_CACHE_TTL = + new ConfigOption<>( + "readiness.cache_ttl", + "How many ms a GET /readiness result is reused before " + + "the storage is probed again, so that several probes " + + "(Kubernetes, load balancers) cost one round of Store " + + "calls per interval; 0 probes on every request.", + rangeInt(0, 60000), + 2000 + ); + + public static final ConfigOption READINESS_MAX_WAITERS = + new ConfigOption<>( + "readiness.max_waiters", + "How many GET /readiness callers may wait at once for " + + "the probe already in flight; the rest get an immediate " + + "503, so a burst of probes during slow storage cannot " + + "hold the REST worker pool.", + rangeInt(1, 10000), + 16 + ); + public static final ConfigOption SERVER_USE_K8S = new ConfigOption<>( "server.use_k8s", diff --git a/hugegraph-server/hugegraph-hbase/src/main/java/org/apache/hugegraph/backend/store/hbase/HbaseSessions.java b/hugegraph-server/hugegraph-hbase/src/main/java/org/apache/hugegraph/backend/store/hbase/HbaseSessions.java index 3255d092d6..e5c72cdda2 100644 --- a/hugegraph-server/hugegraph-hbase/src/main/java/org/apache/hugegraph/backend/store/hbase/HbaseSessions.java +++ b/hugegraph-server/hugegraph-hbase/src/main/java/org/apache/hugegraph/backend/store/hbase/HbaseSessions.java @@ -209,6 +209,18 @@ public void dropNamespace() throws IOException { } } + /** + * Whether the table can serve requests: it exists, is enabled and all its + * regions are assigned (an existing but disabled table is not available). + */ + public boolean tableAvailable(String table) throws IOException { + TableName tableName = TableName.valueOf(this.namespace, table); + try (Admin admin = this.hbase.getAdmin()) { + return admin.tableExists(tableName) && admin.isTableEnabled(tableName) && + admin.isTableAvailable(tableName); + } + } + public boolean existsTable(String table) throws IOException { TableName tableName = TableName.valueOf(this.namespace, table); try (Admin admin = this.hbase.getAdmin()) { diff --git a/hugegraph-server/hugegraph-hbase/src/main/java/org/apache/hugegraph/backend/store/hbase/HbaseStore.java b/hugegraph-server/hugegraph-hbase/src/main/java/org/apache/hugegraph/backend/store/hbase/HbaseStore.java index 56f7210cc8..d68f537ade 100644 --- a/hugegraph-server/hugegraph-hbase/src/main/java/org/apache/hugegraph/backend/store/hbase/HbaseStore.java +++ b/hugegraph-server/hugegraph-hbase/src/main/java/org/apache/hugegraph/backend/store/hbase/HbaseStore.java @@ -20,10 +20,16 @@ import java.io.IOException; import java.util.HashMap; import java.util.Iterator; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.function.BiFunction; import java.util.stream.Collectors; @@ -87,6 +93,83 @@ private void registerMetaHandlers() { HbaseMetrics metrics = new HbaseMetrics(this.sessions); return metrics.compact(this.tableNames()); }); + + this.registerMetaHandler(META_STORAGE_READINESS, (session, meta, args) -> { + E.checkArgument(args.length == 1 && args[0] instanceof Number, + "Expect the timeout in ms as the only argument"); + return this.storageReadiness(((Number) args[0]).longValue()); + }); + } + + public static final String META_STORAGE_READINESS = "storage_readiness"; + + private static final ExecutorService READINESS_EXECUTOR = Executors.newCachedThreadPool(r -> { + Thread thread = new Thread(r, "hbase-readiness"); + thread.setDaemon(true); + return thread; + }); + + /** + * Whether the HBase cluster can serve this graph: admin round trips + * (is every table of the graph enabled and available) within the budget. + * The body carries no addresses and no raw exception text, since the + * readiness endpoint is unauthenticated. + */ + private Map storageReadiness(long timeoutMs) { + List tables = this.tableNames(); + return readinessOf(() -> firstUnavailable(tables, this.sessions::tableAvailable), + timeoutMs, READINESS_EXECUTOR); + } + + /** + * Every table of the graph (vertices, edges, indexes, counters) must be + * enabled and available: the number of the first one that is not, or 0. + * Checked one after another inside the probe's time budget. + */ + public static int firstUnavailable(List tables, TableCheck check) throws Exception { + if (tables.isEmpty()) { + return 1; + } + for (int i = 0; i < tables.size(); i++) { + if (!check.available(tables.get(i))) { + return i + 1; + } + } + return 0; + } + + @FunctionalInterface + public interface TableCheck { + boolean available(String table) throws Exception; + } + + /** The probe outcome of one bounded availability check; public for the unit test. */ + public static Map readinessOf(Callable unavailable, long timeoutMs, + ExecutorService executor) { + E.checkArgument(timeoutMs > 0, "The probe timeout must be > 0, but got %s", timeoutMs); + Map body = new LinkedHashMap<>(); + body.put("ready", false); + long start = System.currentTimeMillis(); + Future check = executor.submit(unavailable); + try { + int missing = check.get(timeoutMs, TimeUnit.MILLISECONDS); + boolean ok = missing == 0; + body.put("ready", ok); + body.put("reason", ok ? "ok" : "table " + missing + " of the graph is not available"); + } catch (TimeoutException e) { + check.cancel(true); + body.put("reason", "hbase did not answer within " + timeoutMs + " ms"); + } catch (ExecutionException e) { + Throwable cause = e.getCause() != null ? e.getCause() : e; + LOG.warn("Storage readiness: the hbase admin call failed", cause); + body.put("reason", "hbase failed: " + cause.getClass().getSimpleName()); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + check.cancel(true); + body.put("reason", "interrupted"); + } + body.put("hbase_millis", System.currentTimeMillis() - start); + return body; } protected void registerTableManager(HugeType type, HbaseTable table) { diff --git a/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreStorageProbe.java b/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreStorageProbe.java new file mode 100644 index 0000000000..c32537913d --- /dev/null +++ b/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreStorageProbe.java @@ -0,0 +1,388 @@ +/* + * 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.store.hstore; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionService; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorCompletionService; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.commons.lang3.concurrent.BasicThreadFactory; +import org.apache.hugegraph.pd.client.PDClient; +import org.apache.hugegraph.pd.common.PDException; +import org.apache.hugegraph.pd.grpc.Metapb; +import org.apache.hugegraph.store.grpc.state.HgStoreStateGrpc; +import org.apache.hugegraph.store.grpc.state.SubStateReq; +import org.apache.hugegraph.util.E; +import org.apache.hugegraph.util.Log; +import org.slf4j.Logger; + +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; +import io.grpc.StatusRuntimeException; + +/** + * Storage-aware readiness of this server, from this server's point of view: + * at least one Store answers a direct, local, read-only gRPC call + * (HgStoreState.getScanState, which reads the node's own scan-pool stats and + * never touches raft). The Store list comes from PD, refreshed in the + * background (single-flight) while the last known list is used right away, so + * PD only matters until the first list is known: a PD that is slow, restarting + * or down afterwards does not change the readiness of a server whose Stores + * still answer. Every known Store is pinged in parallel and the first answer + * wins; every wait is bounded by one shared time budget. The result carries + * no addresses and no raw exception text, since it is served without + * authentication; the full messages go to the log. + *

+ * Scope: ready means this server knows a Store list and reaches at least one + * Store over gRPC. A Store whose status RPC answers while its raft or + * partition path is broken is not detected. A read through the graph path + * was measured as the gate and rejected: a failed read invalidates the + * partition cache the queries share (with PD down the data plane then fails + * too), and one moving partition stalls every read through that client past + * the budget, so a rolling Store restart took every server out of the + * Service while the traffic was fine. + */ +public final class HstoreStorageProbe { + + public static final String META_STORAGE_READINESS = "storage_readiness"; + + private static final Logger LOG = Log.logger(HstoreStorageProbe.class); + + private static final ExecutorService EXECUTOR = Executors.newCachedThreadPool( + new BasicThreadFactory.Builder().namingPattern("storage-readiness-%d") + .daemon(true).build()); + + private static final KnownStores KNOWN = new KnownStores(); + private static final Map CHANNELS = new ConcurrentHashMap<>(); + private static final Set STALE_CHANNELS = ConcurrentHashMap.newKeySet(); + + private HstoreStorageProbe() { + } + + /** The active stores as PD sees them. */ + public interface StoreLister { + + List activeStores() throws Exception; + } + + /** One cheap call to one store; returning (any value) means it answered. */ + public interface StorePinger { + + void ping(Metapb.Store store, long timeoutMs) throws Exception; + } + + /** The last store list PD answered with, shared by consecutive probes. */ + public static final class KnownStores { + + private volatile List stores = Collections.emptyList(); + private volatile long at; + private volatile Boolean pdOk; + private volatile long pdAt; + private final AtomicReference>> inFlight = + new AtomicReference<>(); + + public List stores() { + return this.stores; + } + + public long ageMs() { + return this.at == 0L ? -1L : System.currentTimeMillis() - this.at; + } + + /** Outcome of the last finished PD refresh, null before the first one. */ + public Boolean pdOk() { + return this.pdOk; + } + + public long pdAgeMs() { + return this.pdAt == 0L ? -1L : System.currentTimeMillis() - this.pdAt; + } + + public void update(List stores) { + this.pdOk = true; + this.pdAt = System.currentTimeMillis(); + if (stores != null && !stores.isEmpty()) { + this.stores = Collections.unmodifiableList(new ArrayList<>(stores)); + this.at = this.pdAt; + } + } + + public void pdFailed() { + this.pdOk = false; + this.pdAt = System.currentTimeMillis(); + } + + /** + * The refresh in flight, or a new one started on `executor`: only one + * PD call runs at a time no matter how many probes miss the cache, + * so a hung PD parks one thread, not one per probe. + */ + CompletableFuture> refresh(StoreLister lister, + ExecutorService executor) { + CompletableFuture> running = this.inFlight.get(); + if (running != null && !running.isDone()) { + return running; + } + CompletableFuture> mine = new CompletableFuture<>(); + if (!this.inFlight.compareAndSet(running, mine)) { + return this.inFlight.get(); + } + executor.execute(() -> { + try { + List stores = lister.activeStores(); + this.update(stores); + mine.complete(stores == null ? Collections.emptyList() : stores); + } catch (Throwable e) { + this.pdFailed(); + mine.completeExceptionally(e); + } + }); + return mine; + } + } + + /** The unauthenticated body: no addresses, no raw exception text. */ + private static Map result(boolean ready, String reason, int activeStores, + Long answeredStore, Boolean pdReachable, + long pdAgeMs, long storesAgeMs, long storeMillis) { + Map map = new LinkedHashMap<>(); + map.put("ready", ready); + map.put("reason", reason); + map.put("active_stores", activeStores); + map.put("answered_store", answeredStore); + map.put("pd_reachable", pdReachable); + map.put("pd_checked_age_ms", pdAgeMs); + map.put("stores_age_ms", storesAgeMs); + map.put("store_millis", storeMillis); + return map; + } + + /** + * Probe through the process-wide PD client and this probe's own plaintext + * channels to the stores (the store gRPC server takes no credentials). + * + * @param timeoutMs the whole budget for PD plus stores + */ + public static Map probe(long timeoutMs) { + PDClient pd = HstoreSessionsImpl.getDefaultPdClient(); + if (pd == null) { + return result(false, "pd client not initialised", 0, null, false, -1L, -1L, 0L); + } + return probe(KNOWN, () -> { + List stores = pd.getActiveStores(); + pruneChannels(CHANNELS, stores, STALE_CHANNELS); + return stores; + }, HstoreStorageProbe::pingScanState, timeoutMs, EXECUTOR); + } + + /** + * Channels of stores that left the PD list are shut down only when they + * are absent from two consecutive listings: a probe that snapshotted the + * previous store list may still be pinging them (and a stale-list ping + * through a closed channel would be a false not-ready, cached for the + * TTL). `stale` remembers the addresses missing from the last listing. + */ + static void pruneChannels(Map channels, + List stores, Set stale) { + if (stores == null || stores.isEmpty()) { + return; + } + Set live = new HashSet<>(); + for (Metapb.Store store : stores) { + live.add(store.getAddress()); + } + stale.retainAll(channels.keySet()); + channels.entrySet().removeIf(e -> { + if (live.contains(e.getKey())) { + stale.remove(e.getKey()); + return false; + } + if (stale.add(e.getKey())) { + return false; // first listing without it: keep for the probes in flight + } + e.getValue().shutdownNow(); + stale.remove(e.getKey()); + return true; + }); + } + + private static void pingScanState(Metapb.Store store, long timeoutMs) { + ManagedChannel channel = CHANNELS.computeIfAbsent(store.getAddress(), address -> { + return ManagedChannelBuilder.forTarget(address).usePlaintext().build(); + }); + HgStoreStateGrpc.newBlockingStub(channel) + .withDeadlineAfter(timeoutMs, TimeUnit.MILLISECONDS) + .getScanState(SubStateReq.getDefaultInstance()); + } + + public static Map probe(KnownStores known, StoreLister lister, + StorePinger pinger, long timeoutMs, + ExecutorService executor) { + E.checkArgument(timeoutMs > 0, "The probe timeout must be > 0, but got %s", timeoutMs); + long deadline = System.currentTimeMillis() + timeoutMs; + + // Refresh the store list from PD in the background (single-flight); + // whatever PD answers lands in `known` for this or the next probe + CompletableFuture> refresh = known.refresh(lister, executor); + + List stores = known.stores(); + Boolean pdReachable = null; + if (stores.isEmpty()) { + // Nothing known yet (first probe after start): PD is the only source + try { + stores = await(refresh, deadline); + pdReachable = true; + } catch (TimeoutException e) { + return result(false, "no store list known and pd did not answer within " + + timeoutMs + " ms", 0, null, false, -1L, -1L, 0L); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return result(false, "interrupted while waiting for the store list", + 0, null, false, -1L, -1L, 0L); + } catch (Exception e) { + LOG.warn("Storage readiness: no store list known and pd failed", e); + return result(false, "no store list known and pd failed: " + category(e), + 0, null, false, known.pdAgeMs(), -1L, 0L); + } + if (stores.isEmpty()) { + return result(false, "no active store registered in pd", + 0, null, true, known.pdAgeMs(), known.ageMs(), 0L); + } + } + + long storeStart = System.currentTimeMillis(); + // Ping every known store at once and take the first answer: a store + // whose connection hangs (a pod that just went away) must not eat the + // budget of the stores that are fine, or a rolling restart would + // flap the readiness of every server + CompletionService pings = new ExecutorCompletionService<>(executor); + List> futures = new ArrayList<>(stores.size()); + for (Metapb.Store store : stores) { + futures.add(pings.submit(() -> { + pinger.ping(store, Math.max(1L, deadline - System.currentTimeMillis())); + return store; + })); + } + List failures = new ArrayList<>(); + Map result = null; + try { + for (int done = 0; done < stores.size() && result == null; done++) { + long remaining = deadline - System.currentTimeMillis(); + Future first; + try { + first = remaining > 0 ? + pings.poll(remaining, TimeUnit.MILLISECONDS) : + pings.poll(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + failures.add("interrupted"); + break; + } + if (first == null) { + failures.add((stores.size() - done) + " store(s) did not answer within " + + timeoutMs + " ms"); + break; + } + try { + Metapb.Store store = first.get(); + result = result(true, "ok", stores.size(), store.getId(), + pdState(refresh, known, pdReachable), known.pdAgeMs(), + known.ageMs(), elapsed(storeStart)); + } catch (ExecutionException e) { + Throwable cause = e.getCause() != null ? e.getCause() : e; + LOG.debug("Storage readiness: a store ping failed", cause); + failures.add("a store failed: " + category(cause)); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + failures.add("interrupted"); + break; + } + } + } finally { + for (Future f : futures) { + f.cancel(true); + } + } + if (result != null) { + return result; + } + return result(false, "none of " + stores.size() + " known store(s) answered: " + + String.join("; ", failures), + stores.size(), null, pdState(refresh, known, pdReachable), + known.pdAgeMs(), known.ageMs(), elapsed(storeStart)); + } + + /** + * The outcome of this probe's PD refresh when it already finished, else + * the outcome of the last finished one (null before any finished). + */ + private static Boolean pdState(CompletableFuture refresh, KnownStores known, + Boolean awaited) { + if (awaited != null) { + return awaited; + } + if (refresh.isDone()) { + return !refresh.isCompletedExceptionally(); + } + return known.pdOk(); + } + + private static T await(Future future, long deadline) throws Exception { + long remaining = Math.max(1L, deadline - System.currentTimeMillis()); + try { + return future.get(remaining, TimeUnit.MILLISECONDS); + } catch (ExecutionException e) { + Throwable cause = e.getCause() != null ? e.getCause() : e; + throw cause instanceof Exception ? (Exception) cause : new RuntimeException(cause); + } + } + + private static long elapsed(long since) { + return System.currentTimeMillis() - since; + } + + /** + * A fixed category for the unauthenticated body: the gRPC status code, "pd + * unreachable" or the exception class, never the message (it can carry PD + * peers and Store host names). + */ + static String category(Throwable e) { + if (e instanceof StatusRuntimeException) { + return ((StatusRuntimeException) e).getStatus().getCode().name(); + } + if (e instanceof PDException) { + return "pd unreachable"; + } + return e.getClass().getSimpleName(); + } +} diff --git a/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreStore.java b/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreStore.java index 6439096674..08d87a2aae 100644 --- a/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreStore.java +++ b/hugegraph-server/hugegraph-hstore/src/main/java/org/apache/hugegraph/backend/store/hstore/HstoreStore.java @@ -114,6 +114,11 @@ private void registerMetaHandlers() { HstoreMetrics metrics = new HstoreMetrics(dbsGet.get(), session); return metrics.metrics(); }); + this.registerMetaHandler(HstoreStorageProbe.META_STORAGE_READINESS, (session, meta, args) -> { + E.checkArgument(args.length == 1 && args[0] instanceof Number, + "Expect the timeout in ms as the only argument"); + return HstoreStorageProbe.probe(((Number) args[0]).longValue()); + }); this.registerMetaHandler("mode", (session, meta, args) -> { E.checkArgument(args.length == 1, "The args count of %s must be 1", meta); diff --git a/hugegraph-server/hugegraph-hstore/src/test/java/org/apache/hugegraph/backend/store/hstore/HstoreStorageProbeTest.java b/hugegraph-server/hugegraph-hstore/src/test/java/org/apache/hugegraph/backend/store/hstore/HstoreStorageProbeTest.java new file mode 100644 index 0000000000..46e3bf848c --- /dev/null +++ b/hugegraph-server/hugegraph-hstore/src/test/java/org/apache/hugegraph/backend/store/hstore/HstoreStorageProbeTest.java @@ -0,0 +1,401 @@ +/* + * 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.store.hstore; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; + +import org.apache.hugegraph.backend.store.hstore.HstoreStorageProbe.KnownStores; +import org.apache.hugegraph.pd.common.PDException; +import org.apache.hugegraph.pd.grpc.Metapb; +import org.junit.AfterClass; +import org.junit.Assert; +import org.junit.Test; + +import io.grpc.CallOptions; +import io.grpc.ClientCall; +import io.grpc.ManagedChannel; +import io.grpc.MethodDescriptor; +import io.grpc.Status; + +public class HstoreStorageProbeTest { + + private static final ExecutorService EXECUTOR = Executors.newCachedThreadPool(); + private static final long BUDGET = 500L; + + @AfterClass + public static void shutdown() { + EXECUTOR.shutdownNow(); + } + + private static Metapb.Store store(long id) { + return Metapb.Store.newBuilder().setId(id).setAddress("10.0.0." + id + ":8500") + .setState(Metapb.StoreState.Up).build(); + } + + private static List stores(long... ids) { + return Arrays.stream(ids).mapToObj(HstoreStorageProbeTest::store) + .collect(Collectors.toList()); + } + + private static KnownStores knowing(long... ids) { + KnownStores known = new KnownStores(); + known.update(stores(ids)); + return known; + } + + private static final HstoreStorageProbe.StorePinger ANSWERS = (store, timeout) -> { + }; + + private static String reason(Map body) { + return (String) body.get("reason"); + } + + @Test + public void testReadyWhenPdAndOneStoreAnswer() { + Map r = HstoreStorageProbe.probe(new KnownStores(), () -> stores(1L, 2L, 3L), + ANSWERS, BUDGET, EXECUTOR); + Assert.assertTrue(reason(r), Boolean.TRUE.equals(r.get("ready"))); + Assert.assertEquals(3, r.get("active_stores")); + Assert.assertNotNull(r.get("answered_store")); + Assert.assertEquals(Boolean.TRUE, r.get("pd_reachable")); + Assert.assertEquals("ok", reason(r)); + } + + @Test + public void testFirstAnsweringStoreWinsAfterFailures() { + AtomicInteger pings = new AtomicInteger(); + HstoreStorageProbe.StorePinger onlyThird = (store, timeout) -> { + pings.incrementAndGet(); + if (store.getId() != 3L) { + throw new IllegalStateException("UNAVAILABLE"); + } + }; + Map r = HstoreStorageProbe.probe(knowing(1L, 2L, 3L), () -> stores(1L, 2L, 3L), + onlyThird, BUDGET, EXECUTOR); + Assert.assertTrue(reason(r), Boolean.TRUE.equals(r.get("ready"))); + Assert.assertEquals(3L, r.get("answered_store")); + // the pings run in parallel; the failing ones may or may not have run + Assert.assertTrue(pings.get() >= 1 && pings.get() <= 3); + } + + /** + * The store whose pod just went away hangs until its deadline; it must + * not eat the budget of a store that answers (otherwise a rolling restart + * would flap every server's readiness). + */ + @Test + public void testHungStoreDoesNotHideAnAnsweringOne() { + HstoreStorageProbe.StorePinger onlySecondAnswers = (store, timeout) -> { + if (store.getId() != 2L) { + Thread.sleep(10_000L); + } + }; + long start = System.currentTimeMillis(); + Map r = HstoreStorageProbe.probe(knowing(1L, 2L, 3L), () -> stores(1L, 2L, 3L), + onlySecondAnswers, BUDGET, EXECUTOR); + long took = System.currentTimeMillis() - start; + Assert.assertTrue(reason(r), Boolean.TRUE.equals(r.get("ready"))); + Assert.assertEquals(2L, r.get("answered_store")); + Assert.assertTrue("took " + took, took < BUDGET); + } + + @Test + public void testFirstProbeWithoutKnownStoresNeedsPd() { + Map r = HstoreStorageProbe.probe(new KnownStores(), () -> { + throw new IllegalStateException("UNAVAILABLE: io exception"); + }, ANSWERS, BUDGET, EXECUTOR); + Assert.assertFalse(Boolean.TRUE.equals(r.get("ready"))); + Assert.assertTrue(reason(r), reason(r).startsWith( + "no store list known and pd failed: IllegalStateException")); + Assert.assertEquals(Boolean.FALSE, r.get("pd_reachable")); + } + + @Test + public void testFirstProbeWithHungPdStaysWithinBudget() { + long start = System.currentTimeMillis(); + Map r = HstoreStorageProbe.probe(new KnownStores(), () -> { + Thread.sleep(10_000L); + return stores(1L); + }, ANSWERS, BUDGET, EXECUTOR); + long took = System.currentTimeMillis() - start; + Assert.assertFalse(Boolean.TRUE.equals(r.get("ready"))); + Assert.assertTrue(reason(r), reason(r).contains("pd did not answer within")); + Assert.assertTrue("took " + took, took < BUDGET * 4); + } + + /** + * PD down, restarting or slow must not change the readiness of a server + * whose stores still answer: the last known list is used right away. + */ + @Test + public void testKnownStoresKeepTheServerReadyWhilePdIsDown() { + Map r = HstoreStorageProbe.probe(knowing(1L, 2L), () -> { + throw new IllegalStateException("PD unreachable"); + }, ANSWERS, BUDGET, EXECUTOR); + Assert.assertTrue(reason(r), Boolean.TRUE.equals(r.get("ready"))); + Assert.assertEquals(2, r.get("active_stores")); + } + + @Test + public void testHungPdDoesNotDelayAProbeWithKnownStores() { + long start = System.currentTimeMillis(); + Map r = HstoreStorageProbe.probe(knowing(1L, 2L), () -> { + Thread.sleep(10_000L); + return stores(1L, 2L); + }, ANSWERS, BUDGET, EXECUTOR); + long took = System.currentTimeMillis() - start; + Assert.assertTrue(reason(r), Boolean.TRUE.equals(r.get("ready"))); + // the refresh is still pending, so the outcome of the last finished + // one (the seed) is reported + Assert.assertEquals(Boolean.TRUE, r.get("pd_reachable")); + Assert.assertTrue("took " + took, took < BUDGET); + } + + @Test + public void testLastPdOutcomeIsReportedWhileTheRefreshIsPending() throws Exception { + KnownStores known = knowing(1L); + HstoreStorageProbe.probe(known, () -> { + throw new IllegalStateException("PD unreachable"); + }, ANSWERS, BUDGET, EXECUTOR); + // the seed reports PD ok; wait for the failed refresh to land + for (int i = 0; i < 100 && !Boolean.FALSE.equals(known.pdOk()); i++) { + Thread.sleep(20L); + } + Assert.assertEquals(Boolean.FALSE, known.pdOk()); + Map r = HstoreStorageProbe.probe(known, () -> { + Thread.sleep(10_000L); + return stores(1L); + }, ANSWERS, BUDGET, EXECUTOR); + Assert.assertTrue(Boolean.TRUE.equals(r.get("ready"))); + Assert.assertEquals(Boolean.FALSE, r.get("pd_reachable")); + Assert.assertTrue(r.containsKey("pd_checked_age_ms")); + } + + @Test + public void testPdAnswerUpdatesTheKnownStoresForTheNextProbe() throws Exception { + KnownStores known = knowing(1L); + HstoreStorageProbe.probe(known, () -> stores(1L, 2L, 3L), ANSWERS, BUDGET, EXECUTOR); + for (int i = 0; i < 50 && known.stores().size() != 3; i++) { + Thread.sleep(20L); + } + Assert.assertEquals(3, known.stores().size()); + Assert.assertTrue(known.ageMs() >= 0L); + } + + @Test + public void testEmptyPdAnswerIsNotReadyAndKeepsTheOldList() { + KnownStores fresh = new KnownStores(); + Map r = HstoreStorageProbe.probe(fresh, Collections::emptyList, ANSWERS, + BUDGET, EXECUTOR); + Assert.assertFalse(Boolean.TRUE.equals(r.get("ready"))); + Assert.assertEquals("no active store registered in pd", reason(r)); + KnownStores known = knowing(1L); + HstoreStorageProbe.probe(known, Collections::emptyList, ANSWERS, BUDGET, EXECUTOR); + Assert.assertEquals(1, known.stores().size()); + } + + @Test + public void testNotReadyWhenEveryStoreFails() { + HstoreStorageProbe.StorePinger refused = (store, timeout) -> { + throw new IllegalStateException("connection refused"); + }; + Map r = HstoreStorageProbe.probe(knowing(7L, 8L), () -> stores(7L, 8L), + refused, BUDGET, EXECUTOR); + Assert.assertFalse(Boolean.TRUE.equals(r.get("ready"))); + Assert.assertTrue(reason(r), reason(r).startsWith("none of 2 known store(s) answered")); + Assert.assertTrue(reason(r), reason(r).contains("a store failed: IllegalStateException")); + Assert.assertFalse(reason(r), reason(r).contains("connection refused")); + Assert.assertNull(r.get("answered_store")); + } + + @Test + public void testHungStoresStayWithinTheBudget() { + HstoreStorageProbe.StorePinger hung = (store, timeout) -> { + Thread.sleep(10_000L); + }; + long start = System.currentTimeMillis(); + Map r = HstoreStorageProbe.probe(knowing(1L, 2L, 3L), () -> stores(1L, 2L, 3L), + hung, BUDGET, EXECUTOR); + long took = System.currentTimeMillis() - start; + Assert.assertFalse(Boolean.TRUE.equals(r.get("ready"))); + Assert.assertTrue(reason(r), reason(r).contains("did not answer within")); + Assert.assertTrue("took " + took, took < BUDGET * 4); + } + + @Test + public void testPingGetsTheRemainingBudget() { + List budgets = Collections.synchronizedList(new java.util.ArrayList<>()); + HstoreStorageProbe.probe(knowing(1L), () -> stores(1L), (store, timeout) -> { + budgets.add(timeout); + }, BUDGET, EXECUTOR); + Assert.assertEquals(1, budgets.size()); + Assert.assertTrue(budgets.get(0) > 0L && budgets.get(0) <= BUDGET); + } + + @Test + public void testMapCarriesNoAddresses() { + Map map = HstoreStorageProbe.probe(knowing(1L), () -> stores(1L), + ANSWERS, BUDGET, EXECUTOR); + Assert.assertEquals(true, map.get("ready")); + Assert.assertEquals(1, map.get("active_stores")); + Assert.assertEquals(1L, map.get("answered_store")); + Assert.assertTrue(map.containsKey("pd_reachable")); + Assert.assertTrue(map.containsKey("stores_age_ms")); + Assert.assertFalse(map.toString().contains("10.0.0.")); + } + + @Test + public void testRejectsNonPositiveBudget() { + Assert.assertThrows(IllegalArgumentException.class, () -> { + HstoreStorageProbe.probe(new KnownStores(), Collections::emptyList, ANSWERS, + 0L, EXECUTOR); + }); + } + + /** + * The body is served without authentication: PD peers from the PD client's + * "PD unreachable, pd.peers=..." and Store host names from gRPC's "Unable + * to resolve host ..." must not reach it, only a category. + */ + @Test + public void testReasonCarriesNoPdPeersNorStoreHosts() { + Map pd = HstoreStorageProbe.probe(new KnownStores(), () -> { + throw new PDException(1, "PD unreachable, pd.peers=pd-0.internal:8686,pd-1.internal:8686"); + }, ANSWERS, BUDGET, EXECUTOR); + Assert.assertFalse(Boolean.TRUE.equals(pd.get("ready"))); + Assert.assertEquals("no store list known and pd failed: pd unreachable", reason(pd)); + + HstoreStorageProbe.StorePinger unresolved = (store, timeout) -> { + throw Status.UNAVAILABLE.withDescription( + "Unable to resolve host store-0.hugegraph-store.svc").asRuntimeException(); + }; + Map st = HstoreStorageProbe.probe(knowing(1L), () -> stores(1L), unresolved, + BUDGET, EXECUTOR); + Assert.assertFalse(Boolean.TRUE.equals(st.get("ready"))); + Assert.assertTrue(reason(st), reason(st).contains("a store failed: UNAVAILABLE")); + String all = pd.toString() + st.toString(); + Assert.assertFalse(all, all.contains("internal") || all.contains("svc") || + all.contains("8686")); + } + + /** A hung PD parks one refresh, not one per probe. */ + @Test + public void testRefreshIsSingleFlight() throws Exception { + AtomicInteger calls = new AtomicInteger(); + KnownStores known = knowing(1L); + HstoreStorageProbe.StoreLister hung = () -> { + calls.incrementAndGet(); + Thread.sleep(3_000L); + return stores(1L, 2L); + }; + for (int i = 0; i < 5; i++) { + Assert.assertEquals(Boolean.TRUE, HstoreStorageProbe.probe(known, hung, ANSWERS, + BUDGET, EXECUTOR) + .get("ready")); + } + Thread.sleep(200L); + Assert.assertEquals(1, calls.get()); + for (int i = 0; i < 40 && known.stores().size() != 2; i++) { + Thread.sleep(100L); + } + Assert.assertEquals(2, known.stores().size()); + HstoreStorageProbe.probe(known, hung, ANSWERS, BUDGET, EXECUTOR); + Thread.sleep(100L); + Assert.assertEquals("a finished refresh allows a new one", 2, calls.get()); + } + + private static final class FakeChannel extends ManagedChannel { + + boolean shut; + + @Override + public ManagedChannel shutdown() { + this.shut = true; + return this; + } + + @Override + public boolean isShutdown() { + return this.shut; + } + + @Override + public boolean isTerminated() { + return this.shut; + } + + @Override + public ManagedChannel shutdownNow() { + return this.shutdown(); + } + + @Override + public boolean awaitTermination(long timeout, TimeUnit unit) { + return true; + } + + @Override + public ClientCall newCall(MethodDescriptor method, + CallOptions options) { + throw new UnsupportedOperationException(); + } + + @Override + public String authority() { + return "fake"; + } + } + + @Test + public void testChannelsOfReplacedStoresAreShutDown() { + Map channels = new java.util.concurrent.ConcurrentHashMap<>(); + FakeChannel kept = new FakeChannel(); + FakeChannel gone = new FakeChannel(); + channels.put("10.0.0.1:8500", kept); + channels.put("10.0.0.9:8500", gone); + java.util.Set stale = new java.util.HashSet<>(); + // first listing without the store: kept, a probe may still use the old list + HstoreStorageProbe.pruneChannels(channels, stores(1L, 2L), stale); + Assert.assertEquals(2, channels.size()); + Assert.assertFalse(gone.shut); + Assert.assertTrue(stale.contains("10.0.0.9:8500")); + // the store comes back: forgotten + HstoreStorageProbe.pruneChannels(channels, stores(1L, 2L, 9L), stale); + Assert.assertFalse(stale.contains("10.0.0.9:8500")); + // gone twice in a row: shut down + HstoreStorageProbe.pruneChannels(channels, stores(1L, 2L), stale); + HstoreStorageProbe.pruneChannels(channels, stores(1L, 2L), stale); + Assert.assertEquals(1, channels.size()); + Assert.assertFalse(kept.shut); + Assert.assertTrue(gone.shut); + HstoreStorageProbe.pruneChannels(channels, null, stale); + Assert.assertEquals("a failed listing prunes nothing", 1, channels.size()); + HstoreStorageProbe.pruneChannels(channels, Collections.emptyList(), stale); + Assert.assertEquals("an empty listing keeps the channels the pings still use", + 1, channels.size()); + Assert.assertFalse(kept.shut); + } +} diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/api/ApiTestSuite.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/api/ApiTestSuite.java index f4d74e49b3..a2ebf17926 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/api/ApiTestSuite.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/api/ApiTestSuite.java @@ -35,6 +35,7 @@ TaskApiTest.class, GremlinApiTest.class, MetricsApiTest.class, + ReadinessApiTest.class, UserApiTest.class, LoginApiTest.class, ProjectApiTest.class, diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/api/ReadinessApiTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/api/ReadinessApiTest.java new file mode 100644 index 0000000000..ea08c7451b --- /dev/null +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/api/ReadinessApiTest.java @@ -0,0 +1,81 @@ +/* + * 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.api; + +import java.util.Map; + +import org.apache.hugegraph.testutil.Assert; +import org.apache.hugegraph.util.JsonUtil; +import org.junit.Test; + +import jakarta.ws.rs.client.ClientBuilder; +import jakarta.ws.rs.core.Response; + +/** + * The readiness endpoint answers 200 on a healthy server whatever the + * backend: "embedded" for the in-process backends of the API suite, "hstore" + * (with the Store fields) on the hstore job, "hbase" (with hbase_millis) on + * the hbase job. + */ +public class ReadinessApiTest extends BaseApiTest { + + private static final String PATH = "/readiness"; + + @Test + public void testReadyOnAHealthyServer() { + Response r = client().get(PATH); + String result = assertResponseStatus(200, r); + Map body = JsonUtil.fromJson(result, Map.class); + Assert.assertEquals(true, body.get("ready")); + String storage = String.valueOf(body.get("storage")); + Assert.assertTrue(storage, "embedded".equals(storage) || "hstore".equals(storage) || + "hbase".equals(storage)); + Assert.assertNotNull(body.get("reason")); + if ("hstore".equals(storage)) { + Assert.assertEquals("ok", body.get("reason")); + Assert.assertTrue(((Number) body.get("active_stores")).intValue() >= 1); + Assert.assertNotNull(body.get("answered_store")); + Assert.assertTrue(body.containsKey("cached")); + } + if ("hbase".equals(storage)) { + Assert.assertEquals("ok", body.get("reason")); + Assert.assertTrue(((Number) body.get("hbase_millis")).longValue() >= 0L); + Assert.assertTrue(body.containsKey("cached")); + } + if (!"embedded".equals(storage)) { + Assert.assertNotNull("one entry per probed backend configuration", body.get("probes")); + Assert.assertFalse("no graph name on the unauthenticated endpoint", result.contains("\"graph\"")); + } + } + + /** + * A Kubernetes httpGet probe carries no credential and no graphspace + * prefix, so the endpoint must answer without either. + */ + @Test + public void testReadyWithoutCredentials() { + Response r = ClientBuilder.newClient().target(BASE_URL + PATH) + .request().get(); + try { + Assert.assertEquals(200, r.getStatus()); + Assert.assertContains("\"ready\":true", r.readEntity(String.class)); + } finally { + r.close(); + } + } +} 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 4189c1692d..f52c555bb1 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 @@ -52,6 +52,8 @@ import org.apache.hugegraph.unit.core.GraphSpaceInfoLocaleTest; import org.apache.hugegraph.unit.core.GraphManagerStoresWaitTest; import org.apache.hugegraph.unit.core.MetaManagerClusterTest; +import org.apache.hugegraph.unit.core.HbaseReadinessTest; +import org.apache.hugegraph.unit.core.StorageReadinessTest; import org.apache.hugegraph.unit.core.DirectionsTest; import org.apache.hugegraph.unit.core.ExceptionTest; import org.apache.hugegraph.unit.core.GraphManagerAdminInitTest; @@ -142,6 +144,8 @@ GraphSpaceInfoLocaleTest.class, GraphManagerStoresWaitTest.class, MetaManagerClusterTest.class, + StorageReadinessTest.class, + HbaseReadinessTest.class, DirectionsTest.class, SerialEnumTest.class, diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/api/filter/LoadDetectFilterTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/api/filter/LoadDetectFilterTest.java index 5be5b64a92..1cfd22e757 100644 --- a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/api/filter/LoadDetectFilterTest.java +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/api/filter/LoadDetectFilterTest.java @@ -112,6 +112,23 @@ public void testFilter_WhiteListPathIgnored() { Assert.assertTrue(this.testAppender.events().isEmpty()); } + /** + * A readiness probe must answer from the storage state, not from the + * worker load: a Server that is merely busy is still ready, and shedding + * probes would pull every busy Server out of the Service during a spike. + */ + @Test + public void testFilter_ReadinessIgnoredLikeVersions() { + setupPath("readiness", List.of("readiness")); + this.setConfigProvider(createConfig(2, 0)); + this.workLoad.incrementAndGet(); + + this.loadDetectFilter.filter(this.requestContext); + + Assert.assertEquals(1, this.workLoad.get().get()); + Assert.assertTrue(this.testAppender.events().isEmpty()); + } + @Test public void testFilter_RejectsWhenWorkerLoadIsTooHigh() { setupPath("graphs/hugegraph/vertices", diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/HbaseReadinessTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/HbaseReadinessTest.java new file mode 100644 index 0000000000..fc59de36ec --- /dev/null +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/HbaseReadinessTest.java @@ -0,0 +1,85 @@ +/* + * 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.util.List; +import java.util.Map; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +import org.apache.hugegraph.backend.store.hbase.HbaseStore; +import org.apache.hugegraph.testutil.Assert; +import org.junit.AfterClass; +import org.junit.Test; + +import com.google.common.collect.ImmutableList; + +/** The HBase probe outcome for available tables, one disabled table, a failing and a hung check. */ +public class HbaseReadinessTest { + + private static final ExecutorService EXECUTOR = Executors.newCachedThreadPool(); + + @AfterClass + public static void shutdown() { + EXECUTOR.shutdownNow(); + } + + private static Map readiness(Callable check, long timeout) { + return HbaseStore.readinessOf(check, timeout, EXECUTOR); + } + + @Test + public void testOutcomes() throws Exception { + Map ok = readiness(() -> 0, 500L); + Assert.assertEquals(true, ok.get("ready")); + Assert.assertEquals("ok", ok.get("reason")); + Assert.assertTrue(((Number) ok.get("hbase_millis")).longValue() >= 0L); + + // an existing but disabled table is not available, whichever table it is + List tables = ImmutableList.of("g_v", "g_oe", "g_ie", "g_si"); + Map disabled = readiness(() -> HbaseStore.firstUnavailable(tables, t -> !t.equals("g_ie")), + 500L); + Assert.assertEquals(false, disabled.get("ready")); + Assert.assertEquals("table 3 of the graph is not available", disabled.get("reason")); + Assert.assertEquals(0, HbaseStore.firstUnavailable(tables, t -> true)); + Assert.assertEquals(1, HbaseStore.firstUnavailable(ImmutableList.of(), t -> true)); + // a throwing check surfaces as a failure, not as ready + Assert.assertThrows(java.io.IOException.class, () -> { + HbaseStore.firstUnavailable(tables, t -> { + throw new java.io.IOException("x"); + }); + }); + + Map failed = readiness(() -> { + throw new java.io.IOException("Unable to resolve host hbase-master.svc"); + }, 500L); + Assert.assertEquals(false, failed.get("ready")); + Assert.assertEquals("hbase failed: IOException", failed.get("reason")); + Assert.assertFalse(failed.toString().contains("svc")); + + long start = System.currentTimeMillis(); + Map hung = readiness(() -> { + Thread.sleep(5_000L); + return 0; + }, 300L); + Assert.assertEquals(false, hung.get("ready")); + Assert.assertEquals("hbase did not answer within 300 ms", hung.get("reason")); + Assert.assertTrue(System.currentTimeMillis() - start < 2_000L); + } +} diff --git a/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/StorageReadinessTest.java b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/StorageReadinessTest.java new file mode 100644 index 0000000000..7d6d5e9398 --- /dev/null +++ b/hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/StorageReadinessTest.java @@ -0,0 +1,325 @@ +/* + * 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.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.hugegraph.api.filter.AuthenticationFilter; +import org.apache.hugegraph.api.filter.PathFilter; +import org.apache.hugegraph.api.profile.StorageReadiness; +import org.apache.hugegraph.testutil.Assert; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.mockito.Mockito; + +import jakarta.ws.rs.container.ContainerRequestContext; +import jakarta.ws.rs.core.UriInfo; + +public class StorageReadinessTest { + + @Before + @After + public void reset() { + StorageReadiness.resetCache(); + } + + private static Map result(boolean ready, String reason) { + Map map = new LinkedHashMap<>(); + map.put("ready", ready); + map.put("reason", reason); + map.put("active_stores", 3); + return map; + } + + @Test + public void testReadyBodyAndCacheReuse() { + AtomicInteger probes = new AtomicInteger(); + StorageReadiness.Probe probe = t -> { + probes.incrementAndGet(); + return result(true, "ok"); + }; + Map first = StorageReadiness.check(probe, 1000L, 60_000L); + Map second = StorageReadiness.check(probe, 1000L, 60_000L); + Assert.assertTrue(StorageReadiness.isReady(first)); + Assert.assertEquals("hstore", first.get("storage")); + Assert.assertEquals(false, first.get("cached")); + Assert.assertEquals(true, second.get("cached")); + Assert.assertEquals(1, probes.get()); + } + + @Test + public void testZeroTtlProbesEveryTime() { + AtomicInteger probes = new AtomicInteger(); + StorageReadiness.Probe probe = t -> { + probes.incrementAndGet(); + return result(false, "no active store registered in pd"); + }; + StorageReadiness.check(probe, 1000L, 0L); + Map body = StorageReadiness.check(probe, 1000L, 0L); + Assert.assertFalse(StorageReadiness.isReady(body)); + Assert.assertEquals(2, probes.get()); + Assert.assertEquals("no active store registered in pd", body.get("reason")); + } + + @Test + public void testProbeFailureIsNotReadyWithReason() { + StorageReadiness.Probe probe = t -> { + throw new IllegalStateException("The 'hugegraph' store of hstore has not been opened"); + }; + Map body = StorageReadiness.check(probe, 1000L, 0L); + Assert.assertFalse(StorageReadiness.isReady(body)); + Assert.assertEquals("probe failed: IllegalStateException", body.get("reason")); + // the endpoint is unauthenticated: no raw message in the body + Assert.assertFalse(body.toString().contains("has not been opened")); + } + + @Test + public void testTimeoutIsPassedToTheProbe() { + StorageReadiness.Probe probe = t -> result(true, "budget " + t); + Map body = StorageReadiness.check(probe, 750L, 0L); + Assert.assertEquals("budget 750", body.get("reason")); + } + + @Test + public void testCachedCopyIsIsolated() { + StorageReadiness.Probe probe = t -> result(true, "ok"); + Map first = StorageReadiness.check(probe, 1000L, 60_000L); + first.put("ready", false); + Map second = StorageReadiness.check(probe, 1000L, 60_000L); + Assert.assertTrue(StorageReadiness.isReady(second)); + } + + /** + * Concurrent callers with no cache share one probe: the first runs it on + * its own thread, the others wait for that result. Nothing holds a + * monitor while the probe does its I/O. + */ + @Test + public void testConcurrentCallersShareOneProbe() throws Exception { + AtomicInteger probes = new AtomicInteger(); + CountDownLatch started = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + StorageReadiness.Probe probe = t -> { + probes.incrementAndGet(); + started.countDown(); + release.await(5, TimeUnit.SECONDS); + return result(true, "ok"); + }; + ExecutorService pool = Executors.newFixedThreadPool(4); + try { + Future> owner = pool.submit(() -> { + return StorageReadiness.check(probe, 2000L, 0L); + }); + Assert.assertTrue(started.await(2, TimeUnit.SECONDS)); + List>> followers = new ArrayList<>(); + for (int i = 0; i < 3; i++) { + followers.add(pool.submit(() -> StorageReadiness.check(probe, 2000L, 0L))); + } + Thread.sleep(100L); + release.countDown(); + Assert.assertTrue(StorageReadiness.isReady(owner.get(2, TimeUnit.SECONDS))); + Assert.assertEquals(false, owner.get().get("cached")); + for (Future> f : followers) { + Map body = f.get(2, TimeUnit.SECONDS); + Assert.assertTrue(StorageReadiness.isReady(body)); + Assert.assertEquals(true, body.get("shared")); + } + Assert.assertEquals(1, probes.get()); + } finally { + release.countDown(); + pool.shutdownNow(); + } + } + + /** Beyond maxWaiters, callers get an immediate 503 instead of a worker-pool slot. */ + @Test + public void testExcessWaitersAreRejectedAtOnce() throws Exception { + CountDownLatch started = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + StorageReadiness.Probe probe = t -> { + started.countDown(); + release.await(5, TimeUnit.SECONDS); + return result(true, "ok"); + }; + ExecutorService pool = Executors.newFixedThreadPool(3); + try { + Future> owner = pool.submit(() -> { + return StorageReadiness.check("hstore", probe, 5000L, 0L, 1); + }); + Assert.assertTrue(started.await(2, TimeUnit.SECONDS)); + Future> waiter = pool.submit(() -> { + return StorageReadiness.check("hstore", probe, 5000L, 0L, 1); + }); + Thread.sleep(100L); + long start = System.currentTimeMillis(); + Map rejected = StorageReadiness.check("hstore", probe, 5000L, 0L, 1); + long took = System.currentTimeMillis() - start; + Assert.assertFalse(StorageReadiness.isReady(rejected)); + Assert.assertEquals("too many readiness callers waiting for the probe (1)", + rejected.get("reason")); + Assert.assertTrue("took " + took, took < 500L); + release.countDown(); + Assert.assertTrue(StorageReadiness.isReady(owner.get(2, TimeUnit.SECONDS))); + Map shared = waiter.get(2, TimeUnit.SECONDS); + Assert.assertTrue(StorageReadiness.isReady(shared)); + Assert.assertEquals(true, shared.get("shared")); + // the slot is free again + Assert.assertTrue(StorageReadiness.isReady( + StorageReadiness.check("hstore", t -> result(true, "ok"), 1000L, 0L, 1))); + } finally { + release.countDown(); + pool.shutdownNow(); + } + } + + /** Independent backend configurations are probed side by side; one failing makes the server not ready. */ + @Test + public void testEveryRemoteConfigurationIsProbed() throws Exception { + AtomicInteger healthy = new AtomicInteger(); + AtomicInteger failing = new AtomicInteger(); + List remotes = new ArrayList<>(); + remotes.add(new StorageReadiness.RemoteGraph("g1", "hbase", t -> { + healthy.incrementAndGet(); + return result(true, "ok"); + })); + remotes.add(new StorageReadiness.RemoteGraph("g2", "hbase", t -> { + failing.incrementAndGet(); + return result(false, "hbase failed: IOException"); + })); + Map body = StorageReadiness.probeAll(remotes, 1000L); + Assert.assertEquals(1, healthy.get()); + Assert.assertEquals(1, failing.get()); + Assert.assertFalse(StorageReadiness.isReady(body)); + Assert.assertEquals("hbase configuration 2 of 2: hbase failed: IOException", body.get("reason")); + List probes = (List) body.get("probes"); + Assert.assertEquals(2, probes.size()); + Assert.assertEquals(false, ((Map) probes.get(1)).get("ready")); + // the unauthenticated body carries no graph names + Assert.assertFalse(body.toString().contains("g2")); + Assert.assertNull(((Map) probes.get(0)).get("graph")); + + // a hung configuration is bounded by the shared budget + remotes.add(new StorageReadiness.RemoteGraph("g3", "hstore", t -> { + Thread.sleep(5_000L); + return result(true, "ok"); + })); + long start = System.currentTimeMillis(); + body = StorageReadiness.probeAll(remotes, 300L); + Assert.assertTrue(System.currentTimeMillis() - start < 2_000L); + Assert.assertFalse(StorageReadiness.isReady(body)); + Assert.assertEquals(3, ((List) body.get("probes")).size()); + + // a single configuration keeps the plain body plus its one probe entry + body = StorageReadiness.probeAll(remotes.subList(0, 1), 1000L); + Assert.assertTrue(StorageReadiness.isReady(body)); + Assert.assertEquals(1, ((List) body.get("probes")).size()); + } + + /** Every async probe runs as the internal admin: the auth context is a thread local. */ + @Test + public void testAsyncProbesCarryTheAdminContext() throws Exception { + List remotes = new ArrayList<>(); + List users = java.util.Collections.synchronizedList(new ArrayList<>()); + for (int i = 0; i < 3; i++) { + remotes.add(new StorageReadiness.RemoteGraph("g" + i, "hstore", t -> { + org.apache.hugegraph.auth.HugeGraphAuthProxy.Context ctx = + org.apache.hugegraph.auth.HugeGraphAuthProxy.getContext(); + users.add(ctx == null ? "none" : ctx.user().username()); + return result(true, "ok"); + })); + } + Map body = StorageReadiness.probeAll(remotes, 1000L); + Assert.assertTrue(StorageReadiness.isReady(body)); + Assert.assertEquals(3, users.size()); + for (String u : users) { + Assert.assertEquals("admin", u); + } + } + + /** The storage name in the body is the probed backend's. */ + @Test + public void testStorageNameFollowsTheBackend() { + Map body = StorageReadiness.check("hbase", t -> result(true, "ok"), + 1000L, 0L, 4); + Assert.assertEquals("hbase", body.get("storage")); + StorageReadiness.resetCache(); + body = StorageReadiness.check("hbase", t -> { + throw new IllegalStateException("x"); + }, 1000L, 0L, 4); + Assert.assertEquals("hbase", body.get("storage")); + Assert.assertFalse(StorageReadiness.isReady(body)); + } + + /** A follower waits at most its own timeout for the running probe. */ + @Test + public void testFollowerWaitIsBounded() throws Exception { + CountDownLatch started = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + StorageReadiness.Probe probe = t -> { + started.countDown(); + release.await(5, TimeUnit.SECONDS); + return result(true, "ok"); + }; + ExecutorService pool = Executors.newFixedThreadPool(2); + try { + Future> owner = pool.submit(() -> { + return StorageReadiness.check(probe, 5000L, 0L); + }); + Assert.assertTrue(started.await(2, TimeUnit.SECONDS)); + long start = System.currentTimeMillis(); + Map follower = StorageReadiness.check(probe, 200L, 0L); + long took = System.currentTimeMillis() - start; + Assert.assertFalse(StorageReadiness.isReady(follower)); + Assert.assertEquals("a probe is still running after 200 ms", follower.get("reason")); + Assert.assertTrue("took " + took, took >= 200L && took < 2000L); + release.countDown(); + Assert.assertTrue(StorageReadiness.isReady(owner.get(2, TimeUnit.SECONDS))); + } finally { + release.countDown(); + pool.shutdownNow(); + } + } + + /** + * A Kubernetes httpGet probe carries no credential and no graphspace, so + * the endpoint must pass both the graphspace path rewrite and the auth + * filter, exactly like /versions. + */ + @Test + public void testReadinessBypassesPathAndAuthFilters() { + Assert.assertTrue(PathFilter.isWhiteAPI("readiness")); + Assert.assertTrue(PathFilter.isWhiteAPI("versions")); + UriInfo uri = Mockito.mock(UriInfo.class); + Mockito.when(uri.getPath()).thenReturn("readiness"); + ContainerRequestContext ctx = Mockito.mock(ContainerRequestContext.class); + Mockito.when(ctx.getUriInfo()).thenReturn(uri); + Assert.assertTrue(AuthenticationFilter.isWhiteAPI(ctx)); + Mockito.when(uri.getPath()).thenReturn("readiness/"); + Assert.assertFalse(AuthenticationFilter.isWhiteAPI(ctx)); + } +}