Skip to content

fix(server): converge schema caches across HStore servers - #3237

Open
bitflicker64 wants to merge 9 commits into
apache:masterfrom
hugegraph:fix/schema-cache-convergence-3235
Open

bitflicker64 wants to merge 9 commits into
apache:masterfrom
hugegraph:fix/schema-cache-convergence-3235

Conversation

@bitflicker64

@bitflicker64 bitflicker64 commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Purpose of the PR

With several Servers on HStore, a schema change made through one Server can stay invisible on the others indefinitely. Updates and status flips publish no PD event, on purpose, because one event per updateSchemaStatus() call is a broadcast storm. Events that are published can be lost, because the PD watch starts from "now" after a reconnect (#3157). This PR makes every Server converge within a bounded time without adding any event.

The design and its open questions are in discussion #3205. The questions had no maintainer answer yet, so this PR uses conservative defaults. Each one is listed in #3235 with how to override it.

Main Changes

A schema change writes a per-graph version to PD; every other Server checks it every 10 s and clears its cache once when it changed. A burst of about 5,000 updates cost 4,980 clears per remote Server with notify-on-every-update, and 7 with this PR

  • CachedSchemaTransactionV2 writes a new opaque version of the graph to PD after every committed schema change: add, update, status flip, remove, clear, store clear or truncate. The key is HUGEGRAPH/{cluster}/SCHEMA_VERSION/{graphspace}/{graph}. It sits outside GRAPHSPACE/{gs}/{graph}/SCHEMA because clearAllSchema() deletes that prefix.
  • New SchemaVersionReconciler: one task per open graph on one daemon thread per JVM. Every schema.sync.reconcile_interval seconds it reads the version. If the version changed, it clears that graph's schema cache and adopts the value it read before the clear. Versions are only compared for equality.
  • A failed version write never fails the DDL call. It logs a WARN and is retried on the next tick. Existing events are not changed: adds still publish one event per operation, and a failed event put still fails the add as it does today.
  • Cache guard: a clear caused by a remote change and a reader's cache update are now exclusive on the graph's SchemaCaches. A reader skips its update if such a clear ran after its storage read, so data read before the clear is never cached after it. Cache hit paths are unchanged, except that an id-cache hit takes the lock once when it promotes the element to the array cache.
  • StandardHugeGraph.close() stops the graph's reconciler. Dropping a graph deletes its version key, in StandardHugeGraph.drop() and, for the PD mode drop that calls clearBackend() and close() instead, in GraphManager.dropGraph().
  • New per-graph options: schema.sync.enabled (default true) and schema.sync.reconcile_interval (default 10, range 0 to 3600 seconds, 0 turns the checks off). With schema.sync.enabled=false there is no reconciler and no PD traffic, except one delete of the version key when a graph is dropped in PD mode. The cache guard above stays on either way.

Bounds, from the measurements below:

  • A change becomes visible on every other Server within one interval plus one PD read after the version is written, whatever happens to the watch.
  • After a PD outage, the interval counts from when that Server's own PD client can read again.
  • A burst of M changes lasting D seconds costs at most ceil(D / T) + 1 clears per Server, so cluster-wide clears no longer grow with M. Notifying on every update costs M x (N - 1).

Verifying these changes

  • Need tests and can be verified as follows:
    • SchemaVersionReconcilerTest (16 tests) covers the reconciler with an in-memory store:
      • clears once per change, whatever the number of writes in between;
      • adopts the version read before the clear;
      • does not skip its own writes;
      • retries a failed write, including a failure during the retry;
      • survives read failures;
      • stops for a closed graph;
      • is re-armed after an Error.
    • CachedSchemaTransactionTest has 8 new tests:
      • the real getSchema(type, id) and getAllSchema paths don't cache data across a remote clear injected during the read;
      • a local clear doesn't change the generation;
      • registry start and stop, including a closed graph;
      • the two options;
      • the PD key.
    • mvn test -pl hugegraph-server/hugegraph-test -am -P unit-test on Java 11.0.32 at a4f7ae1c: 774 run, 0 failures, 1 skipped. Master 83ef9f3fa has 748 run, 0 failures, 1 skipped. mvn editorconfig:check is clean on the changed modules.
    • Live, 2026-09-24, on kind: 1 PD, 1 Store, 3 Servers, installed with the chart from feat(helm): add HStore deployment chart hugegraph/hugegraph#221, probing each Server pod directly. Before each update test every schema cache was warmed on all Servers and left for 22 s, so the name-miss reload triggered by the task scheduler could not refresh a cache by accident.
    • "Naive" is master plus notifySchemaCacheClear() in updateSchema(), and an INFO line per remote clear so the clears can be counted.
Change through server-0, seen on server-1 / server-2 master 83ef9f3fa naive this PR (e2701127)
add a property key (event published) 0.05 s / 0.06 s 0.05 s / 0.07 s 0.06 s / 0.08 s
append user_data to it (no event) still old after 60 s 0.04 s / 0.06 s 5.32 s / 5.97 s
append again right after the PD pod restarts still old after 60 s 0.03 s / 0.06 s 9.30 s / 14.40 s
full cache clears per Server during 60 s of updates, counted over 90 s 0 (M = 5,860) 4,980 on each remote (M = 4,980) 7 on each Server (M = 5,000)

The live runs used e2701127. The current head a4f7ae1c adds two fixes from review on hugegraph#236, both in the graph drop path, which none of these measurements touch: drop() reads schema.sync.enabled from the configuration, and the PD mode drop also deletes the version key (GraphManagerDropGraphTest). A second commit, be932423, documents the rollout conditions from review here; it changes only the conf template comment.

The 14.40 s is one skipped tick. Server-2's own PD client had not reconnected yet when its tick ran, about 3 s after the write landed through server-0. The next tick picked the change up. From that Server's log:

11:49:07 [WARN] SchemaVersionReconciler - PD unreachable, schema version reconcile skipped for graph 'DEFAULT-hugegraph': ... Failed to get 'HUGEGRAPH/...
11:49:17 [INFO] SchemaVersionReconciler - PD reachable again, schema version reconcile resumed for graph 'DEFAULT-hugegraph'
11:49:17 [INFO] SchemaVersionReconciler - Schema cache of graph 'DEFAULT-hugegraph' cleared by version reconciler (...-3 -> 1790250543...)

Known limits, also in #3235:

  • If the JVM that made a change dies before its failed version write is retried, that change stays unsignalled until the next schema change of the graph.
  • All graphs of a JVM share one reconciler thread. While PD is unreachable, a slow read delays the other graphs' ticks.
  • A reconciler task killed by an Error is re-armed only when a new schema transaction of the graph is created.
  • Rollout: the bound holds only when every Server of the graph runs this release with schema.sync.enabled=true, schema.sync.reconcile_interval above 0 and PD reachable. During a rolling upgrade an older Server writes no version for schema updates, and upgrading it later does not publish the missed change. After the last Server is upgraded, make one schema change or restart the Servers, so every Server clears its schema cache and reloads it from PD. The conf template and doc(config): add schema.sync options hugegraph-doc#498 say this.
  • The options should match on every Server. A Server with schema.sync.enabled=false writes no version, so its updates are not seen by the others.

Does this PR potentially affect the following parts?

Documentation Status

  • Doc - TODO: required documentation is pending; complete it before merging.
  • Doc - Done: documentation is included here or linked below.
  • Doc - No Need: no user-visible documentation is affected.

Documentation files in this PR or paired hugegraph-doc PR:

  • In this PR: hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties (commented examples).
  • Paired doc PR: doc(config): add schema.sync options hugegraph-doc#498 adds schema.sync.enabled and schema.sync.reconcile_interval to the EN and CN config option pages.

@codecov

codecov Bot commented Sep 24, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 67.96537% with 74 lines in your changes missing coverage. Please review.
✅ Project coverage is 41.49%. Comparing base (83ef9f3) to head (8d6e186).
⚠️ Report is 4 commits behind head on master.

Files with missing lines Patch % Lines
...gegraph/backend/cache/SchemaVersionReconciler.java 60.00% 27 Missing and 7 partials ⚠️
...graph/backend/cache/CachedSchemaTransactionV2.java 73.83% 13 Missing and 15 partials ⚠️
...n/java/org/apache/hugegraph/StandardHugeGraph.java 10.00% 7 Missing and 2 partials ⚠️
...n/java/org/apache/hugegraph/core/GraphManager.java 50.00% 2 Missing ⚠️
...ache/hugegraph/meta/managers/GraphMetaManager.java 87.50% 0 Missing and 1 partial ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##             master    #3237      +/-   ##
============================================
+ Coverage     41.35%   41.49%   +0.14%     
- Complexity     7299     7376      +77     
============================================
  Files           802      804       +2     
  Lines         69688    69991     +303     
  Branches       9291     9337      +46     
============================================
+ Hits          28816    29044     +228     
- Misses        37576    37630      +54     
- Partials       3296     3317      +21     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

A schema change made through one Server could stay invisible on the
others indefinitely: updates and status flips publish no PD event, and
the PD watch does not replay events lost while it reconnects.

Every schema change now writes a per-graph opaque version to PD at
HUGEGRAPH/{cluster}/SCHEMA_VERSION/{graphspace}/{graph}. Each Server
checks it once per schema.sync.reconcile_interval (10s by default) and
clears that graph's schema cache when it changed, so a change becomes
visible everywhere within about one interval without relying on events.
A burst of changes costs each Server at most one clear per interval.
Existing events are unchanged.

A clear caused by a remote change and a reader's cache update are now
exclusive on the graph's SchemaCaches, and a reader skips its update if
such a clear ran after it read, so data read before the clear can't
be cached after it.

New options: schema.sync.enabled (default true) and
schema.sync.reconcile_interval (0 to 3600 seconds, 0 disables checks).

Fixes apache#3235
@bitflicker64
bitflicker64 force-pushed the fix/schema-cache-convergence-3235 branch from ec6f9ed to a4f7ae1 Compare September 24, 2026 13:04

@imbajin imbajin left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocking: no newly verified runtime regression. Score: 8.5/10 after five independent review lanes and synthesis; please clarify the rollout preconditions below before release. This assessment uses code/test inspection and CI evidence, not an independently rerun multi-Server experiment; current CI is not fully green.

The convergence bound needs every Server of the graph on a release with
these options, sync enabled, polling on and PD reachable. Say so in the
conf template, with what to do after a rolling upgrade: older Servers
write no version for schema updates.
bitflicker64 added a commit to hugegraph/hugegraph-doc that referenced this pull request Sep 24, 2026
Per review on apache/hugegraph#3237: the convergence bound needs every
server on a release with these options, sync enabled, polling on and PD
reachable, and a schema change or restart after a rolling upgrade.
imbajin

This comment was marked as outdated.

@imbajin imbajin left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Changes requested: revisit the propagation design

The earlier 9/10 recommendation is withdrawn pending design agreement. The burst benchmark is useful, but does not establish that routine, low-frequency DDL needs continuous polling.

Polling versus notification and ACK recovery

Current PR Alternative to evaluate
Normal update/status change Write a version; peers detect it on their next poll Notify peers; ACK after successful cache invalidation
No Schema changes Every Server checks each open graph periodically No periodic Schema-version reads required by this mechanism
Missed delivery Next successful version poll detects the change PD retries unacknowledged updates, coalesced to the latest version per graph

The current PR still preserves existing add notifications. The alternative is a design proposal, not an implemented guarantee.

Please confirm before implementation

  1. Commit and failover. establish a reliable link between the Schema commit and version update. Recover the latest version and confirmation state after PD leadership changes; duplicate delivery is acceptable, lost updates are not.
  2. Bounded retry. retry asynchronously without blocking DDL. Expire delivery attempts for long-offline instances; distinguish restarted instances from old sessions.
  3. Rejoin before serving. an expired instance must reconcile with PD and invalidate/reload stale Schema before accepting graph requests. Do not assume reconnect already guarantees this.

Please compare the smallest notification/ACK/recovery design against polling under realistic DDL frequency and graph/Server counts. Restore normal update/status notification as part of that evaluation; batch redundant invalidations within one operation where practical.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

An unresolved critical stale-cache race and moderate lifecycle and retry-state issues must be fixed.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 1 High severity

Open (1)
What changed in this PR

This PR adds bounded schema-cache convergence across HStore servers using PD-backed per-graph versions and periodic reconciliation.

Changes:

  • Adds schema version writes, reconciliation, retries, and configuration options.
  • Protects cache updates during remote invalidation.
  • Adds lifecycle cleanup, documentation, and comprehensive tests.
File Summary
hugegraph-server/​hugegraph-test/​src/​main/​java/​org/​apache/​hugegraph/​unit/​UnitTestSuite.java Registers the new tests.
hugegraph-server/​hugegraph-test/​src/​main/​java/​org/​apache/​hugegraph/​unit/​core/​GraphManagerDropGraphTest.java Tests version-key cleanup during graph drops.
hugegraph-server/​hugegraph-test/​src/​main/​java/​org/​apache/​hugegraph/​unit/​cache/​SchemaVersionReconcilerTest.java Tests reconciliation, retries, failures, and scheduling.
hugegraph-server/​hugegraph-test/​src/​main/​java/​org/​apache/​hugegraph/​unit/​cache/​CachedSchemaTransactionTest.java Tests cache guards, options, metadata, and lifecycle behavior.
hugegraph-server/​hugegraph-dist/​src/​assembly/​static/​conf/​graphs/​hugegraph.properties Documents synchronization configuration.
hugegraph-server/​hugegraph-core/​src/​main/​java/​org/​apache/​hugegraph/​StandardHugeGraph.java Manages reconciler shutdown and version-key cleanup.
hugegraph-server/​hugegraph-core/​src/​main/​java/​org/​apache/​hugegraph/​meta/​MetaManager.java Exposes schema-version metadata operations.
hugegraph-server/​hugegraph-core/​src/​main/​java/​org/​apache/​hugegraph/​meta/​managers/​GraphMetaManager.java Implements schema-version key handling.
hugegraph-server/​hugegraph-core/​src/​main/​java/​org/​apache/​hugegraph/​config/​CoreOptions.java Defines synchronization options.
hugegraph-server/​hugegraph-core/​src/​main/​java/​org/​apache/​hugegraph/​backend/​cache/​SchemaVersionReconciler.java Implements periodic PD reconciliation and retries.
hugegraph-server/​hugegraph-core/​src/​main/​java/​org/​apache/​hugegraph/​backend/​cache/​CachedSchemaTransactionV2.java Integrates version bumps and guarded cache updates.
hugegraph-server/​hugegraph-api/​src/​main/​java/​org/​apache/​hugegraph/​core/​GraphManager.java Removes version metadata during graph drops.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

addSchema() and updateSchema() cached the written element without the
generation check the read paths use. A remote clear landing between the
storage write and the cache update could then be undone by an element
another server had already replaced. Capture the generation before the
write and update the caches under the arrayCaches lock only if it is
unchanged.
070a961 skipped the cache update of a local write if the generation
changed during it. The reconciler also clears on this server's own
version writes, so a reader could cache the old element after that clear
and the skipped update left it there: an index label stayed CREATING in
the cache after its rebuild set CREATED, failing two hstore core tests.

Invalidate the element and reset the cached-all flag of its type instead,
so the next read loads it from storage.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🔵 Needs a closer look

One or more issues must be addressed before approval.

Review effort: Lite
Findings: None

Resolved since last review (1)

@imbajin imbajin left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocking: yes. Summary: A schema mutation can be committed without publishing the version that tells peers to invalidate their caches, leaving them on stale schema after a writer crash. Evidence: CachedSchemaTransactionV2 writes the schema before a separate version update; SchemaVersionReconciler.run() skips when the stored version still matches applied.

this.updateCache(schema);
this.updateCache(schema, generation);

this.bumpSchemaVersion();

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ The schema write completes in super.addSchema() before this separate PD version write; updateSchema() and removeSchema() use the same ordering. If the writer exits between those calls, peers keep the previously applied version and SchemaVersionReconciler skips them, so they can continue serving stale schema until another mutation or a peer restart. Please make the mutation notification recoverable across this crash window, for example by atomically storing a revision or publishing a repair revision during graph startup, and add a recovery test.

@bitflicker64 bitflicker64 left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocking: no. Summary: One new Minor finding: the version key deleted on graph drop can be written again before the graph's reconciler stops, so the key the drop path is meant to remove can survive in PD. The design question and the writer-crash window are already open in earlier reviews and are not repeated here. Evidence: static reading of GraphManager.dropGraph(), StandardHugeGraph.drop()/close() and SchemaVersionReconciler.run()/stop() at 7f5672a; latest-head CI checks all pass.

g.clearBackend();
try {
// The schema version kept in PD for the schema cache sync
this.metaManager.deleteSchemaVersion(graphSpace, name);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor: The version key is deleted while the graph's reconciler is still scheduled. SchemaVersionReconciler is stopped only later, in g.close() via stopReconciler(). If any bump() during the drop failed (for example the one from CachedSchemaTransactionV2.clear() inside clearBackend()), pendingWrite is true, and a tick that runs between this delete and g.close() passes the graphClosed check and writes a fresh version at SchemaVersionReconciler.run() line 134. stop() uses future.cancel(false), so a tick already running also completes its write. The key then outlives the graph. StandardHugeGraph.drop() (line 1172) has the same order: delete first, close() last. Please delete the key after the reconciler is stopped, for example by moving the delete after g.close() here and after this.close() in drop(), and extend GraphManagerDropGraphTest to cover a pending write at drop time.

@bitflicker64

Copy link
Copy Markdown
Contributor Author

Design update: push from the commit plus a Server-held lease, replacing the 10 s poll

@imbajin Following your review (O(1) notification, not polling) and your design doc of 2026-09-26, I redid the design from scratch instead of tuning the poll. This comment is the proposed design and how it was reached. Nothing new is pushed to this branch yet; I would like your agreement on the design and on the questions at the end before implementing it here in stages.

Schema sync: DDL commit pushes a frame to every Server; each Server renews a 12 s lease every second, and PD answers only after a ReadIndex, a term check and a flush of earlier frames

How the design was reached

  1. Verified the code first. PD's KV write path, the KV watch, HgPdWatch and Pulse, the Server DDL paths, LockUtil and the schema caches. This turned up the master bugs listed below, which any design has to fix.

  2. Five independent designs from different starting points:

    • your doc as written;
    • one aiming to beat it (every DDL waits for all Servers, as in Chubby);
    • one built on the existing KV watch;
    • one with a Server-held lease and no per-session state in PD;
    • one adapted from prior art: Chubby, etcd watch progress, Kubernetes informers, ZooKeeper, and F1 online schema change.

    Two more specs went deep on single layers: the PD commit layer and the Server cache.

  3. An adversarial review of each of the seven. Every claim was re-checked against the code, and reviewers built failure interleavings. Every design had at least one flaw that lost a change or broke its bound. Each flaw had a local fix.

  4. One combined design from the strongest mechanism of each, then three more review rounds on it: correctness under failure, cost and fit with your doc, and a final pass on what changed.

  5. A prototype of the protocol core on 3 PD, 1 Store and 3 Servers, with fault injection. The numbers are below.

The design

  1. Commit. Each DDL step is one PD raft entry: a new KvService.txn carrying the element keys plus a per-graph record {rev, incarnation, state}. rev is the raft log index of that entry. The TXN compares each element's stored JSON with what the Server read (absent for a create) and the record's incarnation. So:

    • a DDL is all or nothing;
    • two Servers cannot silently overwrite each other;
    • a write for a dropped graph's old incarnation fails.

    A commit timeout resends the same request id; PD deduplicates inside apply and never re-proposes. Record states are CREATING, LIVE, DROPPING and DROPPED.

  2. Notify. The apply thread appends each committed TXN to a bounded in-memory ring and never blocks. One sender on the leader merges per session and per graph, and sends the latest revision plus the keys that changed.

    • It sends at once when idle, then at most one frame per graph per 200 ms.
    • It sends only to Servers that have the graph open.
    • Frames go over one new bidi method, KvService.sync, on its own channel.

    The KV watch and the HgPdWatch subjects keep their current behaviour, so Store is not affected.

  3. ACK and lease. Every second the Server sends an echo stamped with its own monotonic send time t0 and the highest revision it has invalidated. PD answers with Alive only after three steps:

    1. a ReadIndex (a quorum confirms it is still leader);
    2. a check that its term is the one the session registered under;
    3. a flush of every earlier frame on that stream.

    Processing that Alive proves the Server invalidated everything committed before t0, so it may serve until t0 + 12 s. Without a renewal by then, its gate answers 503 with Retry-After. The ACK is cumulative, per session, in memory, with no per-message retry.

  4. Handshake. Hello registers the session with its open graphs, then comes a ReadIndex, then a snapshot of those graphs' records. The Server invalidates only graphs whose revision or incarnation moved, so a PD failover does not flush every cache. A graph opened later subscribes the same way and is registered before its first schema read.

  5. Server cache. One immutable snapshot per graph, tagged with the revision it covers, replaces the id, name and array caches. A request pins one snapshot for its duration. A frame invalidates at once, then reads fetch only the changed keys at a proven revision; any gap falls back to a full load. The writer handles its own frame like every other Server (your 4.3).

  6. Rollout.

    • PD goes first. The TXN op stays off until the leader has confirmed that every PD member supports it, because an old follower skips an unknown raft op.
    • A Server detects an old PD by UNIMPLEMENTED and keeps the current behaviour.
    • Servers roll with DDL frozen.

Staleness bound. If a DDL is acknowledged at time T, then from T + 12 s every Server either serves it or answers 503 on that graph. The bound depends only on the Server's own clock rate. It holds under a Server stall, half-open TCP, a PD leader change including a partitioned old leader, and crashes. Requests already running keep the schema they started with.

Measured

Prototype of the PD TXN with the record, the ring and sender, KvService.sync, and the Server session, lease, gate and TXN writes. It was built on master 83ef9f3fa and is not pushed yet. Run on 2026-09-27 on one 12-core host with 3 PD, 1 Store and 3 Servers on 127.0.0.1, Java 11. Server 0 made every DDL.

Test Result
200 sequential DDLs PD apply to cache clear on the other Servers: p50 1 ms, p99 1 ms. Client ack to visible over REST: p50 8.2 ms, p99 16.8 ms (mostly the reload).
1,000 PK creates in 10 s 50 to 51 frames and clears per Server, no 503 answers, all 1,000 listed on every Server.
10 min idle 600 echoes and 600 Alives per Server; 0 frames, 0 KV reads, 0 KV writes from sync.
PD leader kill -9 stream back in 2.4 s, 0 Servers fenced. Leader transfer: 0.39 s, 0 fenced.
SIGSTOP a Server 30 s during DDL 503 right after resume, fresh reads 80 ms later, 0 stale reads.
Half-open link (a proxy that stops forwarding) lease expired at the last good echo + 12.001 s; fresh 200 190 ms after the link returned.
DDL during a handshake; graph opened mid-session no change lost in either.
Tests PD core 112 of 112; unit 775; hstore core 816 on the second run (the first hit an EdgeCoreTest setup that races an async job against a 100 ms sleep).

One claim failed. A SIGSTOP of the first PD leader after cluster start leaves no new leader for about 11 s, and all 3 Servers fenced for 0.6 to 1.8 s. The cause is jraft's pre-vote, which blocks on the connect timeout while holding the node lock. PD sets that timeout from raft.rpc-timeout = 10 s (RaftEngine.java:136). Later leaders fail over in about 2 s. A separate 1 s connect timeout fixes it (Q7 below).

Idle cost at 20 Servers is 72,000 small echo pairs per hour whatever the number of graphs, against 720,000 PD reads per hour for this PR's current poll with 100 open graphs. A DDL costs 1 raft entry instead of 4 to 10 puts today.

Master bugs found on the way (fixed by this design)

  • Lost updates across Servers. LockUtil is JVM-local:
    • a same-name create on two Servers leaves two ID keys and one NAME key;
    • append and eliminate write JSON built from a possibly stale cache;
    • two IL creates on one base label can drop an IL id from indexLabels.
  • Non-atomic DDL. Every element write is an ID put then a NAME put, and multi-element DDLs span more puts.
  • Cache races in CachedSchemaTransactionV2. Local clears do not bump the generation, so a reader can re-install a deleted label. The getAllSchema hit path walks the cache without the lock and can return a truncated or empty list.
  • Stale PD reads. There is no ReadIndex, so a new leader can serve reads before applying earlier-term entries.
  • The KV watch notifies outside apply. It runs on the gRPC thread after the put, so changes applied by a new leader are never announced and order can invert.
  • Drop/recreate race. Drop deletes GRAPH_CONF before clearing schema, so a late clear can wipe a recreated graph.
  • Truncate resets the id counter. hstore truncate calls resetIdByKey (HstoreSessionsImpl.java:806) while the schema stays in PD, so new schema ids can collide with live ones.
  • Store's schema cache only expires. Store's SchemaDriver splits the clear event on - (SchemaDriver.java:210), but since fix(server): sync hstore schema cache clears #3011 the Server writes a JSON value. Updates never emit the event anyway, so the cache refreshes only when its entries expire.
  • Builders mutate cached objects. Schema builders change the cached object before the PD write succeeds, and copy() is shallow.

Where it differs from your doc

  • The Server fences when its lease ends, not at disconnect. A half-open Server does not know it is disconnected, and fencing at disconnect would 503 every Server at every PD election. Measured: 0 fenced on kill, transfer and the later-leader SIGSTOP.
  • No per-message ACK or retry timers. One HTTP/2 stream delivers in order or fails, and a failed stream resyncs from the durable records. A per-message ACK adds no bound over the cumulative one.
  • The revision is the raft log index. It is monotonic per graph but not dense, and replay after a snapshot is safe with a skip rule.
  • Transport. It is a new method on an existing service rather than the existing watch path, so the watch behaviour Store shares stays unchanged.
  • An immutable per-graph snapshot replaces fine-grained caches under a generation.
  • Changed keys in the frame. This is not the object-level delta sync your section 1 excludes. It keeps no history, replays nothing and ships no values; the Server still invalidates first and reads authoritative values, and falls back to a full load on any gap. Without it, a 1,000-element schema script at 20 Servers moves the storm from clears to full reloads.
  • The echo replaces transport keepalive. PD's gRPC server keeps the library default policy, which rejects client keepalive pings more often than every 5 minutes, and no PD client sets keepalive, so there is none to reuse.
  • Barrier DDL is proposed separately (next section).

Not in this PR: barrier DDL

A wait-for-all-Servers step before index rebuild and label or index removal is designed, but it failed every review round. The last gap is a Store write that the client gave up on but Store still applies after the barrier; only Store-side fencing closes that. Cross-Server rebuild and removal are also unsafe on master today, independent of the cache. So it is a separate proposal, as your TODO 2 frames it. In this PR those DDLs still commit atomically and reach every Server within the 12 s bound, and a Server that has seen DELETING rejects new writes to that label.

Stages on this branch

  • A. Master bug fixes that stand alone: the truncate id reset, ReadIndex in IdMetaStore.getId, and a synchronous snapshot checkpoint.
  • B. Copy-on-write at the schema mutation sites.
  • C. The PD TXN and record, dormant behind the flag.
  • D. Server DDL through the TXN, with compares and migration of existing graphs.
  • E. Sync as one unit: stream, handshake, lease, gate and snapshot cache. The poll, the version key and the event puts are removed.

No stage ships sessions without the lease and the resync that recover them.

Questions

  1. Host sync on KvService next to txn (my recommendation) or as a new method on HgPdWatch?
  2. Defaults: echo every 1 s, lease 12 s, coalescing window 200 ms. Fence at lease end only, as proposed?
  3. Store's schema cache: fix the event parse and have Store watch the record keys in this PR, or in a follow-up?
  4. Store-side fencing and the barrier proposal: take them as a follow-up design, and in which order?
  5. Make GRAPH_CONF and the schema prefixes TXN-only once the flag is on (plain puts refused)?
  6. Keep schema on hstore truncate? The id reset is removed either way.
  7. Give PD a separate raft connect timeout (1 s) so the first-leader failover takes about 2 s instead of 11 s?

The full specification is about 20,000 words; I will put the parts you want into the PR body and the doc PR (apache/hugegraph-doc#498) once the design is agreed. The blocking point in your inline comment (a writer crash between the schema write and the version write) goes away with the atomic TXN; I will answer it there.

@bitflicker64

Copy link
Copy Markdown
Contributor Author

@imbajin A second, independent review of the design above, checked against the code, led to these changes. Nothing else in the design moves.

  1. Read-your-writes on the writer. A request pins one schema snapshot, so a Gremlin script that creates a property key and then a vertex label using it would fail: the builder looks the key up in the pinned snapshot (VertexLabelBuilder.java:403-409). Graph create hits the same path for the task schema (TaskTransaction.java:106-128). The fix: the TXN reply carries prev_revision, and the writer applies its own commit to its cache before the DDL returns. The request then moves to the new snapshot. This also removes a wait of up to 200 ms per DDL for the writer's own frame.
  2. Tasks wait for their Server. A rebuild or remove job can run on a Server that has not yet seen the label it works on. Rebuild then fails with "Undefined index label" (IndexLabelRebuildJob.java:176-189), and some remove jobs do nothing (IndexLabelRemoveJob.java:43-47). The fix: a task reads the graph record with a ReadIndex at start and waits until its Server's cache covers it.
  3. Q5 and Q7 are requirements, not options. Refusing plain writes to schema keys once the flag is on (Q5) and a 1 s raft connect timeout (Q7) are both needed for the 12 s bound.
  4. No old PD binary after enable. An old binary skips every TXN entry. After the flag is on, a Server that finds an old PD fails closed with 503 instead of falling back to legacy mode. Downgrading PD goes through a disable runbook.
  5. Correction to the master bug list. The cache "generation" belongs to this PR branch; master has the same race without any generation.
  6. Correction to the Store item. On a default install the Server's cluster option is hg-test (ServerOptions.java:192), while Store's SchemaDriver hardcodes hg (SchemaDriver.java:67). I have not yet checked on a live cluster what Store then sees. A Store fix (cluster name option, a parse that accepts both event formats, then a record watch) is proposed as a separate PR.
  7. Staging. Stages D and E land together; D is never released without E, so there is no interim polling. Stage E removes only the schema-clear event and the vertex and edge cache events, which no other Server consumes. Graph and graphspace events stay.

Next I start stage A on this branch.

HstoreSession.truncate() called resetIdByKey() for the store's graph
name. For the schema store that is {graphspace}/{graph}/m, so PD dropped
every schema id counter of the graph by prefix, while the schema itself
stays in PD meta and survives the truncate. The truncating Server kept
handing out ids from its cached range, but any other Server, or the same
one after a restart, got ids from 1 again and reused the ids of existing
schema elements.

Stop resetting the counters on truncate.
clearGraphSpace() deletes HUGEGRAPH/{cluster}/GRAPHSPACE/{gs}/AUTH by
prefix before it drops the graphs. Graph keys sit beside AUTH under the
graphspace, so the prefix also matched every key of a graph whose name
starts with AUTH, such as AUTHx, and wiped its schema and config before
the graph's own drop ran.

End the auth prefix with the path delimiter.
IdMetaStore.getId() reads the counter from the local RocksDB and then
writes the next value through raft. The read is outside raft, so a
leader elected a moment ago could read the counter before it had applied
its predecessor's last write and hand out the same id range twice.

getId() now waits for a jraft ReadIndex first: the read index is
confirmed by a quorum in the current term, and the wait ends once the
local applied index reaches it. EAGAIN (no entry of the new term
committed yet) and EBUSY (leadership transfer) are retried 5 times,
20 ms apart; the wait is bounded by raft.rpc-timeout. Stores that are
not replicated skip the wait.
jraft calls onSnapshotSave() on the state machine thread with the
snapshot index set to the last applied index, and applies the next entry
as soon as it returns. PD handed the whole save to a job thread, so the
RocksDB checkpoint could include entries past the snapshot index, and a
node installing that snapshot applied them a second time.

Take the checkpoint before returning; only the compression stays on the
job thread. The checkpoint is hard links under the store's write lock.
A failed checkpoint now completes the snapshot once with an error; it
used to report the error and then continue and report success as well.
PD set jraft's connect, default and install-snapshot rpc timeouts all
from raft.rpc-timeout (10 s). A candidate opens a connection to each
peer while it holds the node lock, and a peer that accepts the TCP
connection but never answers the ping (a stopped process, a lost host)
stalls the election for the whole connect timeout. Losing the host of
the first leader after cluster start left the PD cluster without a
leader for about 11 s.

Add raft.rpc-connect-timeout (default 1000 ms, jraft's own default).
Split out raft.rpc-install-snapshot-timeout as well, so the connect
timeout can be lowered without shortening snapshot installs; its
default 0 keeps the current value, raft.rpc-timeout. raft.rpc-timeout
keeps its meaning for other raft requests.
@bitflicker64

Copy link
Copy Markdown
Contributor Author

Stage A of the design above is pushed: 7f5672a7c..8d6e18684. These are master bugs that stand on their own, so they land before the protocol.

Commit Fix
be867b28c hstore truncate no longer calls resetIdByKey for the schema store, so new schema ids cannot reuse the ids of elements that survived the truncate
c26d96130 the graphspace auth clear deletes .../AUTH/ instead of .../AUTH, which also matched graphs named like AUTHx
2d3dd41a8 IdMetaStore.getId waits for a raft ReadIndex before reading the counter, so a leader elected a moment ago cannot hand out the same id range twice
1d4335fdd onSnapshotSave takes the RocksDB checkpoint on the state machine thread, so a snapshot cannot hold entries past its index; a failed checkpoint now completes once with an error instead of reporting an error and then success
8d6e18684 new raft.rpc-connect-timeout (default 1000 ms), which removes the 11 s election after the first leader's host is lost. raft.rpc-install-snapshot-timeout is split out too, and its default of 0 keeps today's value, raft.rpc-timeout. Both are in hugegraph-pd/docs/configuration.md and the yml templates

Tests on a local 1 PD, 1 Store run, Java 11, 2026-09-27: PD common 83 passed, PD core 118 passed (2 skipped as before), unit 776 passed, hstore core 817 passed. Each commit adds tests that fail with its fix reverted (5 PD core failures, and the truncate and auth tests fail).

What is left on this PR:

  • B: copy-on-write for schema elements. Builders and jobs change cached objects before the PD write succeeds today, and copy() is shallow.
  • C: the PD commit layer, off behind a flag until every PD node supports it:
    • the TXN op with the per-graph record, compares and request-id dedup;
    • the member check before enable;
    • refusing plain writes to schema keys once enabled.
  • D: Server DDL through the TXN, with compares for every DDL kind, migration of existing graphs, system schema, and the CREATING, DROPPING and DROPPED states. D lands together with E, so no release polls.
  • E: KvService.sync, the handshake, the 12 s lease and 503 gate, the per-graph snapshot cache with changed-key reloads, the writer's own-commit wait, the task-start wait, and metrics. E removes the current poll, the version key and the schema-clear event puts.
  • Docs and tests: in-repo docs, a rewrite of doc(config): add schema.sync options hugegraph-doc#498, the fault tests on a 3 PD cluster, and a new PR description.

Separate from this PR: barrier DDL with Store-side fencing (Q4), and the Store schema cache fix (Q3).

@imbajin imbajin left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocking: yes. Summary: The previously reported schema-write/version-publication crash window remains unresolved; this review adds a distinct ReadIndex retry window that can fail ID allocation during leader recovery. Evidence: exact-head static review; latest-head CI/status checks passed; no local tests were run.

*/
private static final long ALIVE_PEERS_REFRESH_MS = 1000L;

private static final int READ_INDEX_RETRIES = 5;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ [P2] Keep retrying transient ReadIndex failures until the configured deadline. waitReadIndex() waits for rpcTimeout (10 seconds by default), but this code schedules only five retries 20 ms apart; after roughly 100 ms, readIndex() completes with UNKNOWN_VALUE. A newly elected JRaft leader returns EAGAIN until an entry in its current term commits, so a quorum that needs longer than 100 ms makes IdMetaStore.getId() fail even while the configured timeout remains. Please retry EAGAIN/EBUSY until the timeoutMs deadline and add coverage for recovery after the current retry window.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Bug] Schema changes on one HStore Server can stay invisible on the others indefinitely

3 participants