Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
abca23b
fix: prevent sender control future race
lokidundun Jul 25, 2026
1f28ea1
fix ci
lokidundun Jul 25, 2026
6672de9
fix: harden sender control lifecycle
lokidundun Jul 26, 2026
3d8441a
Fix sender control future and service cleanup races
lokidundun Jul 26, 2026
187a83f
fix ci
lokidundun Jul 27, 2026
ff76512
delete irrelavant test
lokidundun Jul 27, 2026
61bbf71
chagne the exception
lokidundun Jul 27, 2026
ca692c6
fix: handle sender data path failures
lokidundun Jul 28, 2026
99e825d
fix: complete master future on errors
lokidundun Jul 28, 2026
e60d48d
fix: handle sender failures and simplify tests
lokidundun Jul 29, 2026
f82527a
fix checkstyle
lokidundun Jul 29, 2026
d95a78d
fix: preserve sender and cleanup failures
lokidundun Jul 29, 2026
fb16daf
fix: restrict queued control message API
lokidundun Jul 30, 2026
0d9dea2
fix sender failure propagation and test cleanup
lokidundun Jul 31, 2026
0834ca0
fix ci
lokidundun Jul 31, 2026
a29a85f
fix ci
lokidundun Jul 31, 2026
ff95dfd
fix: worker close-done on registered state and bound integrate-test w…
lokidundun Jul 31, 2026
e280e30
fix ci
lokidundun Jul 31, 2026
570ddc9
fix busy client
lokidundun Aug 1, 2026
4ac16f9
fix ci
lokidundun Aug 1, 2026
c0a98e2
fix: fast fail
lokidundun Aug 1, 2026
af0863a
fix fast fail test
lokidundun Aug 1, 2026
d6361b7
fix ci
lokidundun Aug 1, 2026
5915c93
fix ci
lokidundun Aug 2, 2026
8492fd4
fix ci
lokidundun Aug 2, 2026
5189d79
fix: handle sender failures and partial worker cleanup
lokidundun Sep 18, 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
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,7 @@ public void loadGraph() {
"sending edges", e);
}).join();
this.sendManager.finishSend(MessageType.EDGE);
this.sendManager.checkFatal();
this.sendManager.clearBuffer();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -159,7 +159,7 @@ public synchronized void close() {
LOG.error("Error occurred while closing master service", e);
}

if (!failed && this.bsp4Master != null) {
if (this.inited && !failed && this.bsp4Master != null) {
this.bsp4Master.waitWorkersCloseDone();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,9 +54,8 @@ public void init(Config config) {

@Override
public void close(Config config) {
InetSocketAddress address = this.address();
this.connectionManager.shutdownServer();
LOG.info("DataServerManager closed with address '{}'", address);
LOG.info("DataServerManager closed");

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🧹 This makes DataServerManager safe to close after a partial init, but Managers.closeAll() still stops at the first manager that throws, and the last worker manager, WorkerInputManager, throws on the same path. If initAll() fails before WorkerInputManager.init() runs (port already in use in DataServerManager.init(), or a MinIO error in SnapshotManager.init()), LoadService.fetchers is still all nulls and LoadService.close() NPEs on fetcher.close(). WorkerService.close() then logs Error while closing managers with that NPE, and WorkerInputManager's sendExecutor.shutdown() is skipped.

Probe against this head (real WorkerInputManager with mocked send/snapshot managers, close() without init()):

close threw: java.lang.NullPointerException
  at org.apache.hugegraph.computer.core.worker.load.LoadService.close(LoadService.java:82)
  at org.apache.hugegraph.computer.core.input.WorkerInputManager.close(WorkerInputManager.java:97)

testCloseManagersAfterPartialInitFailure can't see this because last is a mock. Could LoadService.close() skip null fetchers (or Managers.closeAll() close each manager in its own try/catch and rethrow the first failure with the rest suppressed)? Registering a real WorkerInputManager after DataServerManager in that test would cover it.

}

public InetSocketAddress address() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,7 @@ public void startSend(MessageType type) {
.map(this.partitioner::workerId)
.collect(Collectors.toSet());
this.sendControlMessageToWorkers(workerIds, MessageType.START);
this.sender.checkFatal();
LOG.info("Start sending message(type={})", type);
}

Expand All @@ -166,6 +167,7 @@ public void finishSend(MessageType type) {
.map(this.partitioner::workerId)
.collect(Collectors.toSet());
this.sendControlMessageToWorkers(workerIds, MessageType.FINISH);
this.sender.checkFatal();
LOG.info("Finish sending message(type={},count={},bytes={})",
type, stat.messageCount(), stat.messageBytes());
}
Expand All @@ -178,6 +180,10 @@ public void clearBuffer() {
this.buffers.clear();
}

public void checkFatal() {
this.checkException();
}

private void sortIfTargetBufferIsFull(WriteBuffers buffer,
int partitionId,
MessageType type) {
Expand Down Expand Up @@ -286,6 +292,7 @@ private void sendControlMessageToWorkers(Set<Integer> workerIds,
}

private void checkException() {
this.sender.checkFatal();
if (this.exception.get() != null) {
throw new ComputerException("Failed to send message",
this.exception.get());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,4 +45,14 @@ CompletableFuture<Void> send(int workerId, MessageType type)
* an exception is thrown processing message.
*/
void transportExceptionCaught(TransportException cause, ConnectionId connectionId);

/**
* Check whether the sender has encountered a fatal error. Implementations
* that run background threads should propagate the first fatal error to
* callers so that the caller can fail fast instead of hanging on a future
* or barrier.
*/
default void checkFatal() {
// no-op by default
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package org.apache.hugegraph.computer.core.sender;

import java.nio.ByteBuffer;
import java.util.concurrent.CompletableFuture;

import org.apache.hugegraph.computer.core.network.message.MessageType;

Expand All @@ -26,11 +27,18 @@ public class QueuedMessage {
private final int partitionId;
private final MessageType type;
private final ByteBuffer buffer;
private final CompletableFuture<Void> controlFuture;

public QueuedMessage(int partitionId, MessageType type, ByteBuffer buffer) {
this(partitionId, type, buffer, null);
}

QueuedMessage(int partitionId, MessageType type, ByteBuffer buffer,
CompletableFuture<Void> controlFuture) {
this.partitionId = partitionId;
this.type = type;
this.buffer = buffer;
this.controlFuture = controlFuture;
}

public int partitionId() {
Expand All @@ -44,4 +52,8 @@ public MessageType type() {
public ByteBuffer buffer() {
return this.buffer;
}

CompletableFuture<Void> controlFuture() {
return this.controlFuture;
}
}
Loading
Loading