-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathMatchingEngineServer.java
More file actions
71 lines (54 loc) · 2.58 KB
/
Copy pathMatchingEngineServer.java
File metadata and controls
71 lines (54 loc) · 2.58 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
import java.io.IOException;
import java.net.ServerSocket;
import java.net.Socket;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.atomic.AtomicInteger;
public class MatchingEngineServer {
private static final int PORT = 5001;
private static final int EXPECTED_NODES = 2;
public static void main(String[] args) throws IOException, InterruptedException {
// Shared queue between sequence buffer and matching engine
BlockingQueue<Order> queue = new ArrayBlockingQueue<>(200_000);
// Phase 3 — analytics and storage
AnalyticsModule analytics = new AnalyticsModule();
StorageLayer storage = new StorageLayer();
MatchingEngineUI ui = new MatchingEngineUI(); // <--- ADD THIS
ui.setVisible(true);
analytics.setUI(ui);
// Matching engine — same core as Phase 1
MatchingEngine engine = new MatchingEngine(queue);
engine.setAnalytics(analytics);
engine.setStorage(storage);
engine.setUI(ui);
// Phase 3 — sequence buffer ensures in-order delivery to engine
SequenceBuffer sequenceBuffer = new SequenceBuffer(queue);
// Start background threads
analytics.start();
storage.start();
engine.start();
AtomicInteger connected = new AtomicInteger(0);
Thread[] handlers = new Thread[EXPECTED_NODES];
System.out.printf("[Server] Listening on port %d, waiting for %d nodes...%n",
PORT, EXPECTED_NODES);
try (ServerSocket serverSocket = new ServerSocket(PORT)) {
for (int i = 0; i < EXPECTED_NODES; i++) {
Socket clientSocket = serverSocket.accept();
String name = "Node-" + connected.incrementAndGet();
handlers[i] = new ClientHandler(clientSocket, sequenceBuffer, name);
handlers[i].start();
System.out.printf("[Server] %s connected. (%d/%d)%n",
name, connected.get(), EXPECTED_NODES);
}
} // ServerSocket closes here — no new connections accepted
// Wait for all node connections to finish
for (Thread h : handlers) h.join();
// Small buffer for last orders to arrive through sequence buffer
Thread.sleep(300);
// Signal engine to shut down
queue.put(Order.POISON_PILL);
System.out.println("[Server] All nodes done. Poison pill sent.");
engine.join();
System.out.println("[Server] Shutdown complete.");
}
}