diff --git a/HISTORY b/HISTORY index a62ec4a..aff4406 100644 --- a/HISTORY +++ b/HISTORY @@ -1,6 +1,7 @@ Version 1.38 2025-11-27 * bugfixed: loadStorageServersFromTracker correctly with rw option + * throw IOException with msg "connection broken" when in.read() < 0 Version 1.37 2025-11-06 * query file info support two combined flags: diff --git a/pom.xml b/pom.xml index cb8f62b..e291567 100644 --- a/pom.xml +++ b/pom.xml @@ -4,7 +4,7 @@ org.csource fastdfs-client-java - 1.37-SNAPSHOT + 1.38-SNAPSHOT fastdfs-client-java fastdfs client for java jar diff --git a/src/main/java/org/csource/fastdfs/ProtoCommon.java b/src/main/java/org/csource/fastdfs/ProtoCommon.java index 0b06c5e..273d7fb 100644 --- a/src/main/java/org/csource/fastdfs/ProtoCommon.java +++ b/src/main/java/org/csource/fastdfs/ProtoCommon.java @@ -180,9 +180,13 @@ public class ProtoCommon { long pkg_len; header = new byte[FDFS_PROTO_PKG_LEN_SIZE + 2]; - - if ((bytes = in.read(header)) != header.length) { - throw new IOException("recv package size " + bytes + " != " + header.length); + bytes = in.read(header); + if (bytes != header.length) { + if (bytes < 0) { + throw new IOException("connection broken"); + } else { + throw new IOException("recv package size " + bytes + " != " + header.length); + } } if (header[PROTO_HEADER_CMD_INDEX] != expect_cmd) { @@ -222,7 +226,7 @@ public class ProtoCommon { byte[] body = new byte[(int) header.body_len]; int totalBytes = 0; int remainBytes = (int) header.body_len; - int bytes; + int bytes = 0; while (totalBytes < header.body_len) { if ((bytes = in.read(body, totalBytes, remainBytes)) < 0) { @@ -234,7 +238,12 @@ public class ProtoCommon { } if (totalBytes != header.body_len) { - throw new IOException("recv package size " + totalBytes + " != " + header.body_len); + if (totalBytes == 0) { + throw new IOException("connection broken"); + } else { + throw new IOException("connection broken, recv length: " + totalBytes + + ", expect length: " + header.body_len); + } } return new RecvPackageInfo((byte) 0, body); diff --git a/src/main/java/org/csource/fastdfs/TrackerClient.java b/src/main/java/org/csource/fastdfs/TrackerClient.java index 565ee56..83d5ad7 100644 --- a/src/main/java/org/csource/fastdfs/TrackerClient.java +++ b/src/main/java/org/csource/fastdfs/TrackerClient.java @@ -85,13 +85,13 @@ public class TrackerClient { connection = trackerServer.getConnection(); } catch (IOException e) { if (failOver) { - System.err.println("trackerServer get connection error, emsg:" + e.getMessage()); + System.err.println("trackerServer get connection error, " + e.getMessage()); } else { throw e; } } catch (MyException e) { if (failOver) { - System.err.println("trackerServer get connection error, emsg:" + e.getMessage()); + System.err.println("trackerServer get connection error, " + e.getMessage()); } else { throw e; } diff --git a/src/main/java/org/csource/fastdfs/pool/ConnectionFactory.java b/src/main/java/org/csource/fastdfs/pool/ConnectionFactory.java index 0a6b663..9e49dee 100644 --- a/src/main/java/org/csource/fastdfs/pool/ConnectionFactory.java +++ b/src/main/java/org/csource/fastdfs/pool/ConnectionFactory.java @@ -23,7 +23,7 @@ public class ConnectionFactory { sock.connect(socketAddress, ClientGlobal.g_connect_timeout); return new Connection(sock, socketAddress); } catch (Exception e) { - throw new MyException("connect to server " + socketAddress.getAddress().getHostAddress() + ":" + socketAddress.getPort() + " fail, emsg:" + e.getMessage()); + throw new MyException("connect to server " + socketAddress.getAddress().getHostAddress() + ":" + socketAddress.getPort() + " fail, emsg: " + e.getMessage()); } } } diff --git a/src/main/java/org/csource/fastdfs/pool/ConnectionManager.java b/src/main/java/org/csource/fastdfs/pool/ConnectionManager.java index fa76798..4f39352 100644 --- a/src/main/java/org/csource/fastdfs/pool/ConnectionManager.java +++ b/src/main/java/org/csource/fastdfs/pool/ConnectionManager.java @@ -62,7 +62,7 @@ public class ConnectionManager { try { isActive = connection.activeTest(); } catch (IOException e) { - System.err.println("send to server[" + inetSocketAddress.getAddress().getHostAddress() + ":" + inetSocketAddress.getPort() + "] active test error ,emsg:" + e.getMessage()); + System.err.println("send to server[" + inetSocketAddress.getAddress().getHostAddress() + ":" + inetSocketAddress.getPort() + "] active test error, emsg: " + e.getMessage()); isActive = false; } if (!isActive) { @@ -84,7 +84,7 @@ public class ConnectionManager { throw new MyException("connect to server " + inetSocketAddress.getAddress().getHostAddress() + ":" + inetSocketAddress.getPort() + " fail, wait_time > " + ClientGlobal.g_connection_pool_max_wait_time_in_ms + "ms"); } catch (InterruptedException e) { e.printStackTrace(); - throw new MyException("connect to server " + inetSocketAddress.getAddress().getHostAddress() + ":" + inetSocketAddress.getPort() + " fail, emsg:" + e.getMessage()); + throw new MyException("connect to server " + inetSocketAddress.getAddress().getHostAddress() + ":" + inetSocketAddress.getPort() + " fail, emsg: " + e.getMessage()); } } return connection; @@ -117,7 +117,7 @@ public class ConnectionManager { connection.closeDirectly(); } } catch (IOException e) { - System.err.println("close socket[" + inetSocketAddress.getAddress().getHostAddress() + ":" + inetSocketAddress.getPort() + "] error ,emsg:" + e.getMessage()); + System.err.println("close socket[" + inetSocketAddress.getAddress().getHostAddress() + ":" + inetSocketAddress.getPort() + "] error, emsg: " + e.getMessage()); e.printStackTrace(); } }