Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
7d1472b
fix: guard transaction and shutdown lifecycles
imbajin Oct 3, 2026
4e30d6d
fix: preserve transaction ownership during cleanup
imbajin Oct 3, 2026
ffe1e54
fix(ci): align Store coverage contract
imbajin Oct 3, 2026
b780352
fix(store): preserve scan cancellation semantics
imbajin Oct 3, 2026
976ecc2
fix(ci): retain native shutdown failure output
imbajin Oct 3, 2026
8989c37
chore: sync lifecycle branch with master
imbajin Oct 4, 2026
de46aee
docs: clarify code and Markdown wrapping
imbajin Oct 4, 2026
9a44b2b
fix(store): retain query cleanup failures
imbajin Oct 4, 2026
91fd925
chore(ci): cancel superseded PR runs
imbajin Oct 4, 2026
c9b928b
chore(ci): cancel superseded PR runs
imbajin Oct 4, 2026
811fc77
chore: sync master into lifecycle PR
imbajin Oct 4, 2026
5267823
fix(store): scope executable JAR to runtime
imbajin Oct 4, 2026
d8c9495
fix(store): drain ordinary scan resources
imbajin Oct 4, 2026
ee6d495
fix(store): preserve query half-close semantics
imbajin Oct 4, 2026
e72b5fc
fix(store): throttle scan cleanup warnings
imbajin Oct 4, 2026
a49d701
fix(store): report cancelled one-shot scans
imbajin Oct 4, 2026
80dd322
fix(store): stabilize shutdown wait checks
imbajin Oct 4, 2026
9ed2884
Merge branch 'apache:master' into master
imbajin Oct 5, 2026
254381b
chore: sync master into lifecycle
imbajin Oct 5, 2026
719222b
fix(store): cancel queries after request half-close
imbajin Oct 5, 2026
de24938
fix(store): synchronize cleanup failure regression
imbajin Oct 5, 2026
e3d9540
fix(rocksdb): dispose final session native owners
imbajin Oct 5, 2026
0cba1c8
fix(test): execute cleanup and OLAP regressions
imbajin Oct 5, 2026
b41b113
fix(cache): preserve graph cache lifetime
imbajin Oct 5, 2026
b7035f8
fix(test): await query transport completion
imbajin Oct 6, 2026
c47eaaa
fix(client): finish drained streams without polling
imbajin Oct 6, 2026
1110c91
fix(cache): finish graph and fixture cleanup
imbajin Oct 6, 2026
f81dc95
fix: preserve cache and scan cleanup lifetimes
imbajin Oct 7, 2026
cf66842
fix: sync lifecycle cleanup with ASF master
imbajin Oct 8, 2026
139c80b
fix(store): exclude unused web test dependencies
imbajin Oct 8, 2026
a73b864
Merge branch 'master' into task/topling-split-lifecycle
imbajin Oct 9, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 23 additions & 0 deletions .github/workflows/pd-store-ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,21 @@ permissions:
contents: read

jobs:
store-stop-shell-test:
permissions:
contents: read
runs-on: ubuntu-24.04
steps:
- name: Checkout
uses: actions/checkout@v7
with:
persist-credentials: false

- name: Verify Store stop failure propagation and PID retention
run: |
docker run --rm -v "$GITHUB_WORKSPACE:/work:ro" -w /work ubuntu:24.04 \
bash hugegraph-server/hugegraph-dist/src/assembly/travis/test-stop-hugegraph-store.sh

rocksdb-compatibility:
permissions:
contents: read
Expand Down Expand Up @@ -431,6 +446,11 @@ jobs:
mvn test -pl hugegraph-store/hg-store-test -am \
-P store-core-test -Djacoco.sessionId=store-core-test

- name: Run server shutdown test
run: |
mvn test -pl hugegraph-store/hg-store-test -am \
-P store-server-test -Djacoco.sessionId=store-server-test

- name: Generate aggregate coverage report
run: |
mvn verify -pl hugegraph-store/hg-store-test -am -P jacoco \
Expand All @@ -450,6 +470,8 @@ jobs:
--require-test-report \
"$TEST_REPORT_DIR/TEST-org.apache.hugegraph.store.core.CoreSuiteTest.xml" \
--require-test-report \
"$TEST_REPORT_DIR/TEST-org.apache.hugegraph.store.service.ServerSuiteTest.xml" \
--require-test-report \
"hugegraph-store/hg-store-node/target/surefire-reports/TEST-org.apache.hugegraph.store.business.StoredRowIngressTest.xml" \
--require-covered-group hg-store-common \
--require-covered-group hg-store-client \
Expand All @@ -460,6 +482,7 @@ jobs:
--require-session store-rocksdb-test \
--require-session store-raftcore-test \
--require-session store-core-test \
--require-session store-server-test \
--require-session store-node-test \
"$REPORT_FILE" \
hg-store-grpc hg-store-common hg-store-client hg-store-rocksdb hg-store-core
Expand Down
3 changes: 0 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -316,9 +316,6 @@ gremlin> :> g.V().limit(5)
```

For comprehensive documentation, visit the [HugeGraph Documentation](https://hugegraph.apache.org/docs/).

See [standalone RocksDB snapshot recovery](docs/rocksdb-recovery.md) before restoring data or mounting store directories.

For an existing deployment, read the [RocksDB upgrade guidance](docs/rocksdb-upgrade.md) before upgrading the storage runtime.

</details>
Expand Down
30 changes: 30 additions & 0 deletions docs/transaction-lifecycle.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
# Request transactions and Store shutdown

Server releases the current thread's graph transactions when a REST request
finishes and when an authenticated context task exits. Cleanup visits every
registered graph even when one graph fails to close, and reports failures.
Cleanup explicitly rolls back unfinished writes, including when `onClose(COMMIT)`
was configured, then releases backend transactions and resets thread-local
transaction behavior and listeners. Applications must commit successful writes
before leaving the request/task boundary. Closing an OLTP traverser preserves
the caller's transaction so that the caller can still commit or roll it back.

The cleanup also applies when schema caches have become cold. Deleting an auth
project or user must remove its associated access or belong edges while
preserving unrelated relationships. The special OLAP vertex path remains
separate; negative schema IDs alone do not identify edges that can be skipped.

Store shutdown refuses new RPCs, cancels active RPCs and scans, and waits for their callbacks and scan/TTL workers before Spring destroys the Store engine and its databases. A failed aggregate-query response callback is logged without skipping cancellation or cleanup waits for other queries. Cancellation may fail in-flight requests; stop writes and check their outcomes before planned maintenance. This is not a guarantee that every in-flight request drains successfully or that a leader transfers without a failover interval.

If a callback cannot finish, or an ordinary scan or bidirectional aggregate query fails to release its plan or partition iterators, shutdown remains pending rather than closing its database underneath it. This also includes iterators discarded while advancing past empty partitions, initializing a sequential scan, or counting rows. RocksDB iterators retain the first failure during automatic or explicit close so that later cleanup cannot hide it; concurrent close calls wait for the release attempt to finish. All scan entry points, including one-shot scans, retain failed cleanup independently of RPC termination and executor shutdown. Failed iterator releases are attempted once and remain diagnostic shutdown blockers. Bidirectional aggregate queries report cleanup failures to the client and log that shutdown is blocked; their final success batch is sent only after cleanup succeeds. The distribution stop script waits up to 30 seconds, returns a nonzero status on timeout, and retains the PID file for diagnosis. Inspect logs and thread dumps before retrying. Do not add a concurrent shutdown hook that closes the same databases.

If shutdown cancellation wins while a one-shot scan releases its iterator, the
response terminates with `CANCELLED`; it does not report successful completion
without its result. Iterator cleanup still finishes before the scan unregisters.

A normal aggregate-query request half-close ends feedback without cancelling already permitted work. A subsequent transport cancellation or deadline still interrupts the workers and releases their resources. A batch-scan RPC accepts one initial query; repeated query requests are ignored before allocating another iterator, including after its final batch. If the remaining feedback credit cannot finish the query, the server returns an explicit query error. Scan task rejection reports `UNAVAILABLE` during shutdown and `RESOURCE_EXHAUSTED` when the running scan pool is full.

See the [Store shutdown instructions](../hugegraph-store/README.md#stopping-a-store-node)
and the Server [module test guidance](../hugegraph-server/AGENTS.md#tests).

Shared schema and element caches retain their invalidation listeners until the graph closes. Request cleanup releases backend leases while preserving those graph caches and their schema identity.
Original file line number Diff line number Diff line change
Expand Up @@ -185,7 +185,7 @@ private static void registerPrivateActions() {
"this$0");
Reflection.registerFieldsToFilter(HugeGraphAuthProxy.Context.class, "ADMIN", "user");
Reflection.registerFieldsToFilter(HugeGraphAuthProxy.ContextTask.class, "runner",
"context");
"cleanup", "context");
Reflection.registerFieldsToFilter(StandardHugeGraph.class, "LOG", "started", "closed",
"mode", "variables", "name", "params", "configuration",
"schemaEventHub", "graphEventHub", "indexEventHub",
Expand All @@ -203,7 +203,8 @@ private static void registerPrivateActions() {
"access$14", "access$15", "access$16", "access$17",
"access$18", "serializer", "loadSchemaStore",
"loadSystemStore", "loadGraphStore", "closeTx",
"analyzer", "serverInfoManager", "reloadRamtable",
"closeCurrentThreadTransaction", "analyzer",
"serverInfoManager", "reloadRamtable",
"reloadRamtable", "access$19", "access$20", "access$21");
Reflection.registerFieldsToFilter(
loadClass("org.apache.hugegraph.StandardHugeGraph$StandardHugeGraphParams"),
Expand Down Expand Up @@ -298,8 +299,10 @@ private static void registerPrivateActions() {
"autoCommit", "beforeRead", "afterWrite", "afterRead",
"commitMutation2Backend", "checkOwnerThread", "doAction",
"store", "reset");
Reflection.registerFieldsToFilter(HugeFactory.class, "LOG", "NAME_REGEX", "graphs");
Reflection.registerMethodsToFilter(HugeFactory.class, "lambda$0");
Reflection.registerFieldsToFilter(HugeFactory.class, "LOG", "NAME_REGEX", "graphs",
"GRAPHS");
Reflection.registerMethodsToFilter(HugeFactory.class, "lambda$0",
"closeCurrentThreadTransactions");
Reflection.registerFieldsToFilter(SchemaElement.class, "graph", "id", "name", "userdata",
"status");
Reflection.registerFieldsToFilter(HugeVertex.class, "EMPTY_SET", "id", "label", "edges",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1291,11 +1291,17 @@ public User user() {
static class ContextTask implements Runnable {

private final Runnable runner;
private final Runnable cleanup;
private final Context context;

public ContextTask(Runnable runner) {
this(runner, HugeFactory::closeCurrentThreadTransactions);
}

ContextTask(Runnable runner, Runnable cleanup) {
this.context = getContext();
this.runner = runner;
this.cleanup = cleanup;
}

@Override
Expand All @@ -1305,7 +1311,7 @@ public void run() {
this.runner.run();
} finally {
try {
HugeFactory.closeCurrentThreadTransactions();
this.cleanup.run();
} catch (Throwable e) {
LOG.error("Failed to close Gremlin worker transactions", e);
} finally {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import java.io.File;
import java.net.URL;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
Expand Down Expand Up @@ -118,6 +119,11 @@ public static void closeCurrentThreadTransactions() {
synchronized (HugeFactory.class) {
graphs = new ArrayList<>(GRAPHS.values());
}
closeCurrentThreadTransactions(graphs);
}

static void closeCurrentThreadTransactions(
Collection<? extends StandardHugeGraph> graphs) {
Throwable failure = null;
for (StandardHugeGraph graph : graphs) {
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1361,22 +1361,25 @@ public void rollback() {
}

public void close() {
try {
this.graphTx.close();
} catch (Exception e) {
LOG.error("Failed to close GraphTransaction", e);
Throwable failure = null;
for (Runnable close : new Runnable[]{this.graphTx::close,
this.systemTx::close,
this.schemaTx::close}) {
try {
close.run();
} catch (RuntimeException | Error error) {
if (failure == null) {
failure = error;
} else if (failure != error) {
failure.addSuppressed(error);
}
}
}

try {
this.systemTx.close();
} catch (Exception e) {
LOG.error("Failed to close SystemTransaction", e);
if (failure instanceof Error) {
throw (Error) failure;
}

try {
this.schemaTx.close();
} catch (Exception e) {
LOG.error("Failed to close SchemaTransaction", e);
if (failure != null) {
throw (RuntimeException) failure;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,14 +56,17 @@ public HugeConfig config() {
public final BackendSession getOrNewSession() {
BackendSession session = this.threadLocalSession.get();
if (session == null) {
session = this.newSession();
assert session != null;
this.threadLocalSession.set(session);
assert !this.sessions.containsKey(Thread.currentThread().getId());
this.sessions.put(Thread.currentThread().getId(), session);
int sessionCount = this.sessionCount.incrementAndGet();
LOG.debug("Now(after connect({})) session count is: {}",
this, sessionCount);
// Serialize new borrowers with the last-session native close.
synchronized (this) {
session = this.newSession();
assert session != null;
this.threadLocalSession.set(session);
assert !this.sessions.containsKey(Thread.currentThread().getId());
this.sessions.put(Thread.currentThread().getId(), session);
int sessionCount = this.sessionCount.incrementAndGet();
LOG.debug("Now(after connect({})) session count is: {}",
this, sessionCount);
}
} else {
this.detectSession(session);
}
Expand Down Expand Up @@ -131,7 +134,7 @@ public void forceResetSessions() {
}
}

public boolean close() {
public synchronized boolean close() {
Pair<Integer, Integer> result = Pair.of(-1, -1);
try {
result = this.closeSession();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -765,7 +765,7 @@ public Iterator<Vertex> queryTaskInfos(Query query) {
}

public Iterator<Vertex> queryTaskInfos(Object... vertexIds) {
if (this.graph().backendStoreFeatures().supportsTaskAndServerVertex()) {
if (this.storeFeatures().supportsTaskAndServerVertex()) {
return this.queryVerticesByIds(vertexIds, false, false,
HugeType.TASK);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -596,7 +596,7 @@ private <V> Iterator<HugeTask<V>> queryTask(Map<String, Object> conditions,
boolean withResult) {
return this.call(() -> {
ConditionQuery query;
if (this.graph.backendStoreFeatures().supportsTaskAndServerVertex()) {
if (this.tx().storeFeatures().supportsTaskAndServerVertex()) {
query = new ConditionQuery(HugeType.TASK);
} else {
query = new ConditionQuery(HugeType.VERTEX);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ protected OltpTraverser(HugeGraph graph) {

@Override
public void close() {
// pass
// The graph's thread-local transaction belongs to the caller.
}

public static void destroy() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ cleanup() {
local status=$?
trap - EXIT
if [[ "$SERVER_START_ATTEMPTED" == "true" ]]; then
"$SERVER_DIR/bin/stop-hugegraph.sh" -m false >/dev/null 2>&1 || status=1
"$SERVER_DIR/bin/stop-hugegraph.sh" -m false || status=1
fi
rm -rf "$WORK_DIR"
exit "$status"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -430,8 +430,8 @@ def validation_command(job):

def reports_for_option(job, option):
pattern = re.escape(option) + (
r'\s+\\?\s*"\$TEST_REPORT_DIR/'
r'(TEST-[A-Za-z0-9_.]+SuiteTest[.]xml)"'
r'\s+\\?\s*"[^"]*/'
r'(TEST-[A-Za-z0-9_.]+Test[.]xml)"'
)
return set(re.findall(pattern, validation_command(job)))

Expand Down Expand Up @@ -492,9 +492,12 @@ assert_order(store_job, [
"-P store-client-test -Djacoco.sessionId=store-client-test",
"-P store-rocksdb-test -Djacoco.sessionId=store-rocksdb-test",
"-P store-raftcore-test -Djacoco.sessionId=store-raftcore-test",
"-P store-core-test -Djacoco.sessionId=store-core-test",
"-P store-server-test -Djacoco.sessionId=store-server-test",
"mvn verify", "--require-session store-common-test",
"--require-session store-client-test", "--require-session store-rocksdb-test",
"--require-session store-raftcore-test", "codecov/codecov-action",
"--require-session store-raftcore-test", "--require-session store-core-test",
"--require-session store-server-test", "codecov/codecov-action",
])
assert store_job.count("mvn clean") == 1
assert "hugegraph-store/hg-store-test/target/site/jacoco/jacoco.xml" in store_job
Expand All @@ -504,16 +507,19 @@ assert "mvn verify -pl hugegraph-store/hg-store-test -am -P jacoco \\ " \
"-DskipTests -Deditorconfig.skip=true -ntp" in " ".join(store_job.split())
assert selected_profiles(store_job, "store") == {
"store-common-test", "store-client-test", "store-rocksdb-test",
"store-raftcore-test", "store-core-test",
"store-raftcore-test", "store-core-test", "store-server-test",
}
assert reports_for_option(store_job, "--require-test-report") == {
"TEST-org.apache.hugegraph.store.common.CommonSuiteTest.xml",
"TEST-org.apache.hugegraph.store.client.ClientSuiteTest.xml",
"TEST-org.apache.hugegraph.store.rocksdb.RocksDbSuiteTest.xml",
"TEST-org.apache.hugegraph.store.raftcore.RaftSuiteTest.xml",
"TEST-org.apache.hugegraph.store.core.CoreSuiteTest.xml",
"TEST-org.apache.hugegraph.store.service.ServerSuiteTest.xml",
"TEST-org.apache.hugegraph.store.business.StoredRowIngressTest.xml",
}
assert not reports_for_option(store_job, "--require-suite-report")
assert values_for_option(store_job, "--require-session") == (selected_profiles(store_job, "store") | {"store-node-test"})
assert values_for_option(store_job, "--require-covered-group") == {
"hg-store-common", "hg-store-client", "hg-store-rocksdb", "hg-store-core",
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
#!/bin/bash
#
# 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.
#
# Verify stop failure propagation and PID retention without signalling real processes.
set -euo pipefail

ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/../../../../.." && pwd)"
SCRIPT="$ROOT/hugegraph-store/hg-store-dist/src/assembly/static/bin/stop-hugegraph-store.sh"
FIXTURE=$(mktemp -d)
trap 'rm -rf "$FIXTURE"' EXIT
mkdir -p "$FIXTURE/bin"
cp "$SCRIPT" "$FIXTURE/bin/stop-hugegraph-store.sh"
cat > "$FIXTURE/bin/util.sh" <<'UTIL'
kill_process_and_wait() {
[[ "$1" == HugeGraphStoreServer && "$2" == 12345 && "$3" == 30 ]] || return 99
return "$STOP_RESULT"
}
UTIL
printf '12345\n' > "$FIXTURE/bin/pid"
if STOP_RESULT=1 bash "$FIXTURE/bin/stop-hugegraph-store.sh"; then
echo "Stop timeout incorrectly returned success" >&2
exit 1
fi
[[ "$(cat "$FIXTURE/bin/pid")" == 12345 ]]
STOP_RESULT=0 bash "$FIXTURE/bin/stop-hugegraph-store.sh"
[[ ! -e "$FIXTURE/bin/pid" ]]
STOP_RESULT=0 bash "$FIXTURE/bin/stop-hugegraph-store.sh"
echo "store-stop-contract-ok"
Loading
Loading