Skip to content

Improve thread safety of shared mutable state #6932

Description

@xxo1shine

Summary

There are several instances of shared mutable state in the P2P / net modules of java-tron that are accessed concurrently by multiple threads, including net workers, synchronization threads, block-fetch thread pools, and scheduled tasks, but are not protected by synchronization, volatile, or thread-safe containers.

Such unsynchronized concurrent reads and writes constitute data races under the Java Memory Model and may result in stale reads, lost updates, inconsistent collection state, and race windows in compound operations such as check-then-set. Under high connection pressure or intensive multi-threaded scheduling, these issues may affect the determinism and stability of node behavior.

We recommend adding appropriate synchronization, visibility, and atomicity guarantees to these shared states.

Root Cause

The following eight instances of shared state lack adequate concurrency protection:

  1. SyncService.syncBlockInProcess uses a non-thread-safe HashSet and is accessed concurrently by multiple threads without synchronization. HashSet does not guarantee internal consistency under concurrent structural modifications.
  2. SyncService.syncNext performs a check-then-set operation on peer.syncChainRequested (checking whether it is unset before assigning a value). The operation is not atomic, leaving a race window when multiple threads enter the code concurrently. This may result in duplicate requests or inconsistent state.
  3. The shared field fetchBlockInfo is not declared volatile and is concurrently read and written by three thread pools. Without a visibility guarantee, a write performed by one thread is not guaranteed to become visible to other threads in a timely manner.
  4. cheatWitnessInfoMap uses a non-thread-safe HashMap. Concurrent reads/writes or structural modifications are not safely supported and may result in lost updates or inconsistent observations. If iteration occurs concurrently with modification, it may also trigger ConcurrentModificationException.
  5. PeerManager.check() modifies the peers collection and related counters without synchronization. Reads, writes, and iteration may therefore occur concurrently and interfere with each other.
  6. BlockChainMetricManager.getDupWitness() has a write-read ordering race involving dupWitnessBlockNum: the block-production thread first calls counterInc (making the counter key visible) and then performs put. Meanwhile, the metrics thread checks the counter key and calls dupWitnessBlockNum.get(witness). If the read occurs between these two operations, it may return null.
  7. MessageCount uses ordinary fields for szCount[], index, and totalCount. add() performs non-atomic read-modify-write operations, while update() updates the rolling window without synchronization. These methods can be invoked concurrently by multiple sender threads, potentially resulting in lost increments and inconsistent counters.
  8. BackupManager.status is an ordinary field updated by the Backup keep-alive scheduled task and UDP event-handling thread, while being read by the consensus block-production thread, Relay scheduled task, and Metrics thread. The field is neither declared volatile nor protected by consistent synchronization. Under the Java Memory Model, reader threads may continue observing a stale role.

The underlying cause is consistent across all cases: shared mutable state is accessed concurrently without adequate synchronization, visibility, or atomicity guarantees.

Impact

  • Concurrent access to shared state may result in stale reads, lost updates, inconsistent collection state, or race conditions in compound operations, affecting the determinism and stability of node behavior.
  • The issues are generally intermittent and dependent on thread scheduling, making them more likely to occur under high connection pressure or intensive multi-threaded workloads.
  • The impact is limited to the runtime stability of an individual node. These issues do not alter protocol semantics and do not affect consensus correctness or asset security.
  • A stale read of BackupManager.status may cause a witness node to incorrectly attempt or skip a scheduled production slot and may temporarily make Relay and monitoring behavior inconsistent.

Suggested Fix

Apply appropriate concurrency protection to each shared state based on its access pattern:

  1. syncBlockInProcess: Replace it with a thread-safe collection, such as Collections.synchronizedSet or ConcurrentHashMap.newKeySet(), or consistently synchronize access to the set.
  2. syncChainRequested in SyncService.syncNext: Replace the check-then-set operation with an atomic operation, such as performing both operations inside a synchronized block or using an atomic reference with CAS, eliminating the race window.
  3. fetchBlockInfo: Declare the field as volatile (or use an atomic reference) to guarantee cross-thread visibility. If compound updates are involved, synchronization should be added as well.
  4. cheatWitnessInfoMap: Replace HashMap with ConcurrentHashMap.
  5. PeerManager.check(): Synchronize reads and writes to peers and the related counters, or use thread-safe containers to ensure that modifications and iteration do not occur concurrently.
  6. getDupWitness / dupWitnessBlockNum: Ensure that the companion map becomes visible before the counter key by changing the write order to dupWitnessBlockNum.put(...) followed by counterInc(...). On the read side, use dupWitnessBlockNum.getOrDefault(witness, 0L) to eliminate the potential NullPointerException caused by automatic unboxing, providing an additional safeguard.
  7. MessageCount: Use a thread-safe implementation and synchronize add, add(int), getCount, and update consistently to ensure cross-thread visibility and prevent lost updates.
  8. BackupManager.status: At minimum, declare the field volatile so that consensus, Relay, Metrics, and other reader threads observe the latest role.

General Principle

For each shared state, identify the threads that access it and the corresponding access patterns, then choose the minimal and correct concurrency mechanism:

  • Visibility → volatile
  • Compound atomicity → synchronization / locks or CAS
  • Concurrent collections → thread-safe implementations

The goal is to avoid unprotected shared mutable state while keeping synchronization overhead and implementation complexity to a minimum.

Activity

  1. jakamobiii commented on Aug 31, 2026

    @jakamobiii

    @xxo1shine For BackupManager.status, could you clarify the call path from a stale read to actual block-production behavior? Does the consensus thread directly rely on this field when deciding whether to participate in a production slot?

  2. xxo1shine commented on Aug 31, 2026

    @xxo1shine
    CollaboratorAuthor

    @jakamobiii Yes. BackupManager.status is read from the block-production path to determine whether the local node currently has the active role required for production. Since the status is updated by the backup keep-alive and UDP handling threads, a stale read may cause the production thread to temporarily act on an outdated role. This does not change consensus validation rules, because an invalid block would still be rejected by the protocol, but it may affect the availability of the individual witness by causing it to attempt or skip production based on stale local state.

  3. halibobo1205 commented on Sep 9, 2026

    @halibobo1205
    Collaborator

    [DISCUSS] A related concurrency issue was observed in BackupServer during CI.

    org.tron.common.backup.BackupServerTest > test FAILED
        org.junit.runners.model.TestTimedOutException: test timed out after 60 seconds
            at sun.misc.Unsafe.park(Native Method)
            at java.util.concurrent.locks.LockSupport.parkNanos(LockSupport.java:215)
            at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.awaitNanos(AbstractQueuedSynchronizer.java:2078)
            at java.util.concurrent.ThreadPoolExecutor.awaitTermination(ThreadPoolExecutor.java:1475)
            at java.util.concurrent.Executors$DelegatedExecutorService.awaitTermination(Executors.java:675)
            at org.tron.common.es.ExecutorServiceManager.shutdownAndAwaitTermination(ExecutorServiceManager.java:90)
            at org.tron.common.backup.socket.BackupServer.close(BackupServer.java:106)
            at org.tron.common.backup.BackupServerTest.tearDown(BackupServerTest.java:42)
    

    If close() runs before bind() publishes channel, it skips closing the channel. A subsequent successful bind leaves the server thread blocked in closeFuture().sync(), causing BackupServerTest to time out during teardown.

    This requires lifecycle coordination: channel can legitimately still be null when shutdown starts, so adding volatile alone would not fix the race.

    This is another concurrency issue in the backup subsystem. It does not address the BackupManager.status visibility issue described in item 8.

  4. lxcmyf commented on Sep 9, 2026

    @lxcmyf
    Collaborator

    One point to tighten for item 3: volatile only solves visibility, not the lifecycle race. fetchBlock() performs a null-check followed by assignment, while fetchBlockProcess(fetchBlockInfo) operates on a previously captured object and later clears this.fetchBlockInfo unconditionally. If the old request is completed and a new request is installed between those steps, the scheduled worker can erase the newer request; concurrent fetch callers can also both observe null and overwrite each other. An AtomicReference lifecycle looks safer: use compareAndSet(null, newInfo) when installing, and compareAndSet(observedInfo, null) on success or timeout. A deterministic test using an old scheduled snapshot followed by a newly installed request would help pin down this case.

  5. xxo1shine commented on Sep 11, 2026

    @xxo1shine
    CollaboratorAuthor

    @halibobo1205, the BackupServer lifecycle race you identified is addressed in PR #6961, which has not yet been merged. After bind completes, the code checks shutdown and immediately closes the newly bound channel if shutdown has already started. The PR also adds a deterministic regression test.

    This is separate from the cross-thread visibility issue involving BackupManager.status in item 8. Declaring status as volatile addresses the stale-read issue described there. The bind/close race involves operation ordering and requires additional lifecycle coordination; volatile alone is insufficient for that race.

  6. xxo1shine commented on Sep 11, 2026

    @xxo1shine
    CollaboratorAuthor

    @lxcmyf, volatile addresses the cross-thread visibility issue originally described in item 3. Here, fetchBlockInfo supports an auxiliary block-fetch optimization: the normal block-fetch request has already been sent before this field is set, so losing this auxiliary state does not mean the original request is lost.

    The compound-operation races you identified are indeed not resolved by volatile, but they have not yet been shown to cause persistent block-fetch failures. Given the scope of this fix, I would prefer to use volatile first and avoid expanding the changes to request lifecycle management. Atomic installation and identity-checked cleanup can be evaluated separately as a follow-up improvement.

  7. waynercheung commented on Sep 11, 2026

    @waynercheung
    Collaborator

    Following the visibility/lifecycle distinction discussed above, I think item 5 could use a more specific invariant: each peer should be removed, accounted for, and cleaned up only once.

    PeerManager already uses a synchronized list and atomic counters, but check() does not acquire the class monitor used by remove(Channel). Both removal paths also decrement the counter without checking whether peers.remove(peer) succeeded.

    One possible interleaving is:

    1. check() takes a snapshot containing a peer whose disconnectTime is older than DISCONNECTION_TIME_OUT.
    2. The regular disconnect callback (P2pEventHandlerImpl.onDisconnect -> PeerManager.remove(Channel)) removes that peer, decrements its counter, and calls onDisconnect().
    3. check() resumes with its old snapshot. Removal now returns false, but it still decrements the counter and calls onDisconnect() again.

    The reverse order (callback locates the peer, check() removes it first) has the same effect. A late callback is exactly the situation check() exists for, so the two paths are expected to be active at the same time by design. The visible effect today is limited to the active/passive counters in node info drifting, since the second onDisconnect() runs on already-cleared state; the point is more that the invariant is not expressed anywhere than that the drift is harmful.

    Could both paths share a removal helper that updates membership and counters under the same lock and reports whether it actually removed the peer? Only the successful caller would perform cleanup, which could remain outside the membership lock.

    That helper is also what makes this testable: check() is private static with no seam to pause after the snapshot, and a PeerConnection built in a unit test has no wired services for onDisconnect(). A regression test at the helper level could remove the same peer twice (once through remove(Channel), once directly) and verify, for both active and passive peers, that the second call reports no removal and the counter is decremented once.

  8. xxo1shine commented on Sep 14, 2026

    @xxo1shine
    CollaboratorAuthor

    @waynercheung, the invariant you identified—each peer should be removed, accounted for, and cleaned up only once—is more precise. check() and the disconnect callback can process the same peer. Currently, both paths decrement the counter even when removal fails, which can cause the counters to drift and onDisconnect() to run twice. I agree that both paths should share a synchronized helper that removes the peer, updates the counter, and reports whether removal succeeded. Only the caller that successfully removes the peer should perform cleanup. The regression test will cover repeated removal of both active and passive peers and verify that the second attempt neither changes the counter nor triggers cleanup.

  9. waynercheung commented on Sep 14, 2026

    @waynercheung
    Collaborator

    Thanks @xxo1shine, that matches what I had in mind for item 5; happy to review the PR.

    For item 2, I traced a possible consequence that may help with prioritization. The proposed atomic check-and-set already addresses the underlying issue.

    SyncService.processBlock() on the peer's message-handling thread and processSyncBlock() on the sync-handle-block thread can both reach syncNext(peer). They do not share a lock protecting the syncChainRequested == null check, and forkLock is acquired only afterward, so both callers can pass the check and send a SyncBlockChainMessage.

    If both requests receive replies, one problematic interleaving is that the first ChainInventoryMessage clears the pending request and the second arrives before another request is registered. ChainInventoryMsgHandler.check() then throws BAD_MESSAGE, which maps to BAD_PROTOCOL and channel.close(BAD_PEER_BAN_TIME), with a one-hour ban duration. This gives a code path from duplicate local requests to penalizing a healthy peer, although I have not reproduced it end to end. The window requires the message-handling caller to be stalled between the check and the assignment (for example, waiting on forkLock during a fork switch), so I would expect this to be rare.

    A dedicated per-peer lock covering the check, summary generation, and assignment looks reasonable for preventing duplicate installation. It would also avoid sharing the peer monitor with checkAndPutAdvInvRequest(). The implementation should account for the block-processing caller already holding blockLock when checking lock order.

    For regression coverage, we could hold the first caller inside summary generation, ensure the second caller has reached lock contention, then release the first and assert that only one SyncBlockChainMessage is sent. The response-side rejection is already covered by ChainInventoryMsgHandlerTest, so the new test only needs the sender side.

  10. xxo1shine commented on Sep 16, 2026

    @xxo1shine
    CollaboratorAuthor

    @waynercheung, thank you for the additional context. The syncNext() call in processBlock() requests the next batch early after a block is received, while the call in processSyncBlock() is triggered only after block processing completes and the fetch queue is empty. They serve different purposes and occur at different stages. Normally, once the former sets syncChainRequested, the latter returns immediately.

    There is no need to introduce a dedicated per-peer lock here. We will reuse the existing lock to keep the null check, state update, and send operation in syncNext() atomic. This closes the duplicate-request window under an extreme interleaving without introducing another lock or additional lock-order complexity.

  11. waynercheung commented on Sep 16, 2026

    @waynercheung
    Collaborator

    Reusing an existing lock works for me; the requirement is to serialize competing syncNext() calls across the null check, assignment, and send.

    One lock-order detail is worth keeping in mind, since #6970 now calls syncNext() from ChainInventoryMsgHandler while holding blockLock. Extending the existing forkLock section looks like a reasonable option, preserving the existing blockLock -> forkLock order on that path. Using blockLock would also follow that order, although it would make message-handling callers contend with block processing.

    The SyncService monitor is the one to avoid here: handleSyncBlock() holds it while acquiring blockLock, so making syncNext() synchronized would introduce the reverse order through ChainInventoryMsgHandler.

    The atomic section also needs to cover the direct caller in SyncService.processBlock(). A same-peer concurrency regression against the real syncNext(), verifying that only one request is sent while it remains outstanding, would confirm that. Happy to review the update.

  12. xxo1shine commented on Sep 17, 2026

    @xxo1shine
    CollaboratorAuthor

    @waynercheung, noted on the lock ordering. We will account for it in the implementation and add the corresponding concurrency regression tests.

  13. removed this from the GreatVoyage-v4.8.3 milestone on Sep 24, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    • Status
      No status

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions