{"record":{"id":"6d516edcd99caace","repo":"apache/seatunnel","slug":"rejected-edge-collector-from-another-collector","errorCode":null,"errorMessage":"Rejected edge collector from {}: another collector is already connected","messagePattern":"Rejected edge collector from (.+?): another collector is already connected","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-edge-socket/src/main/java/org/apache/seatunnel/connectors/seatunnel/edgesocket/source/EdgeSocketIngressServer.java","lineNumber":199,"sourceCode":"            }\n        }\n        log.warn(\"Edge socket receiver loop exception, retrying\", e);\n        return isInterruptedDuringRetryWait();\n    }\n\n    private void serveOneCollector() throws IOException {\n        ServerSocket ss = serverSocket;\n        if (ss == null || ss.isClosed()) {\n            return;\n        }\n\n        Socket raw = ss.accept();\n        raw.setSoTimeout(config.getAcceptTimeoutMs());\n        try (IngressChannel channel = new BlockingIngressChannel(raw)) {\n            log.info(\"Accepted edge collector connection from {}\", channel.remoteAddress());\n            if (hasActiveCollector) {\n                log.warn(\n                        \"Rejected edge collector from {}: another collector is already connected\",\n                        channel.remoteAddress());\n                channel.writeLine(EdgeSocketResponseCode.REJECTED.getCode());\n                return;\n            }\n            activateAndServe(channel);\n        }\n    }\n\n    private void activateAndServe(IngressChannel channel) throws IOException {\n        hasActiveCollector = true;\n        try {\n            if (!protocolHandler.authenticate(channel)) {\n                return;\n            }\n            closeServerSocket();\n            protocolHandler.receiveLoop(channel, this::isReceiverActive);\n        } finally {\n            hasActiveCollector = false;","sourceCodeStart":181,"sourceCodeEnd":217,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-edge-socket/src/main/java/org/apache/seatunnel/connectors/seatunnel/edgesocket/source/EdgeSocketIngressServer.java#L181-L217","documentation":"Warned by serveOneCollector() when a new collector connection arrives while hasActiveCollector is true. The edge-socket source allows only one collector connection at a time, so the new one is refused with the REJECTED response code and closed. This enforces single-writer semantics on the ingress port.","triggerScenarios":"A second TCP connection is accepted while another collector is being served — e.g. duplicate collector instance, aggressive reconnect after a network blip while the old connection lingers, or a monitoring probe connecting to the same port.","commonSituations":"Two collector processes (or replicas) pointed at the same edge-socket port, collector's reconnect backoff shorter than the server's detection of a dead connection, ops port-scanning the ingress port.","solutions":["Ensure exactly one collector instance targets this edge-socket source port.","Add reconnect backoff/jitter on the collector so it does not hammer the port while the previous connection is still considered active.","If the old collector is truly dead, fix/disable keepalive/idle detection so hasActiveCollector clears promptly, or restart the source to free the slot.","Dedicate a separate port for health checks instead of reusing the ingress port."],"exampleFix":"// before (collector)\nwhile (!connected) { reconnect(); } // tight reconnect loop, second connection\n// after\nwhile (!connected) { reconnect(); Thread.sleep(backoffMs); backoffMs = Math.min(backoffMs * 2, maxMs); }","handlingStrategy":"retry","validationCode":"// collector-side check before connecting\nif (anotherCollectorProcessHoldsPort(host, port)) { throw new IllegalStateException(\"edge-socket port already served by another collector\"); }","typeGuard":null,"tryCatchPattern":"String resp = channel.readLine(); if (EdgeSocketResponseCode.REJECTED.getCode().equals(resp)) { scheduleReconnectWithBackoff(); }","preventionTips":["Deploy exactly one collector per edge-socket source/port","Use exponential backoff with jitter on reconnect","Separate health-check traffic from the ingress port","Configure server-side idle timeouts so dead collectors free the single slot quickly"],"tags":["tcp","concurrency","connection-rejected"],"backgroundTag":"connection-refused","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}