From 94774b6c8558c607c158a796d0cf9a4cf6421308 Mon Sep 17 00:00:00 2001 From: re20052 <51275902044@stu.ecnu.edu.cn> Date: Fri, 5 Jun 2026 16:52:25 +0800 Subject: [PATCH] [fix](stream_load) Fix stream load IPv6 host parsing --- .../apache/doris/httpv2/rest/LoadAction.java | 31 ++++---- .../doris/httpv2/rest/LoadActionTest.java | 71 +++++++++++++++++++ 2 files changed, 89 insertions(+), 13 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/LoadAction.java b/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/LoadAction.java index 5e058fcba8fc0f..e4b2e21c8d7532 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/LoadAction.java +++ b/fe/fe-core/src/main/java/org/apache/doris/httpv2/rest/LoadAction.java @@ -46,6 +46,7 @@ import org.apache.doris.thrift.TNetworkAddress; import com.google.common.base.Strings; +import com.google.common.net.HostAndPort; import io.netty.handler.codec.http.HttpHeaderNames; import jakarta.servlet.http.HttpServletRequest; import jakarta.servlet.http.HttpServletResponse; @@ -504,14 +505,11 @@ private TNetworkAddress selectEndpointByRedirectPolicy(HttpServletRequest req, B } String reqHost = ""; - String[] pair = reqHostStr.split(":"); - if (pair.length == 1) { - reqHost = pair[0]; - } else if (pair.length == 2) { - reqHost = pair[0]; - } else { + try { + reqHost = HostAndPort.fromString(reqHostStr).getHost(); + } catch (IllegalArgumentException e) { LOG.info("Invalid header host: {}", reqHostStr); - throw new LoadException("Invalid header host: " + reqHost); + throw new LoadException("Invalid header host: " + reqHostStr); } // User specified redirect policy @@ -579,19 +577,26 @@ private Pair splitHostAndPort(String hostPort) throws AnalysisE throw new AnalysisException("empty endpoint: " + hostPort); } - String[] pair = hostPort.split(":"); - if (pair.length != 2) { + String host; + int port; + try { + HostAndPort hp = HostAndPort.fromString(hostPort); + if (!hp.hasPort()) { + throw new IllegalArgumentException("No port found"); + } + host = hp.getHost(); + port = hp.getPort(); + } catch (IllegalArgumentException e) { LOG.info("Invalid endpoint: {}", hostPort); throw new AnalysisException("Invalid endpoint: " + hostPort); } - int port = Integer.parseInt(pair[1]); if (port <= 0 || port >= 65536) { - LOG.info("Invalid endpoint port: {}", pair[1]); - throw new AnalysisException("Invalid endpoint port: " + pair[1]); + LOG.info("Invalid endpoint port: {}", port); + throw new AnalysisException("Invalid endpoint port: " + port); } - return Pair.of(pair[0], port); + return Pair.of(host, port); } // NOTE: This function can only be used for AuditlogPlugin stream load for now. diff --git a/fe/fe-core/src/test/java/org/apache/doris/httpv2/rest/LoadActionTest.java b/fe/fe-core/src/test/java/org/apache/doris/httpv2/rest/LoadActionTest.java index 4d57c25fd71f31..2c8c59d413a786 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/httpv2/rest/LoadActionTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/httpv2/rest/LoadActionTest.java @@ -352,6 +352,61 @@ public void testRedirectToStreamLoadForwardBuildsForwardUrl() throws Exception { redirectView.getUrl()); } + @Test + public void testSplitHostAndPortParsesIpv4() throws Exception { + LoadAction loadAction = new LoadAction(); + org.apache.doris.common.Pair result = invokeSplitHostAndPort(loadAction, "10.0.0.1:8040"); + Assertions.assertEquals("10.0.0.1", result.first); + Assertions.assertEquals(8040, result.second.intValue()); + } + + @Test + public void testSplitHostAndPortParsesIpv6() throws Exception { + LoadAction loadAction = new LoadAction(); + org.apache.doris.common.Pair result = + invokeSplitHostAndPort(loadAction, "[2001:db8::1]:8040"); + Assertions.assertEquals("2001:db8::1", result.first); + Assertions.assertEquals(8040, result.second.intValue()); + } + + @Test + public void testSelectEndpointByRedirectPolicyParsesIpv4HostAndEndpoint() throws Exception { + LoadAction loadAction = new LoadAction(); + HttpServletRequest request = Mockito.mock(HttpServletRequest.class); + Mockito.when(request.getHeader(LoadAction.HEADER_REDIRECT_POLICY)) + .thenReturn(LoadAction.REDIRECT_POLICY_PRIVATE); + Mockito.when(request.getHeader("host")).thenReturn("10.0.0.1:8030"); + + Backend backend = Mockito.mock(Backend.class); + Mockito.when(backend.getPrivateEndpoint()).thenReturn("192.168.1.1:8040"); + Mockito.when(backend.getPublicEndpoint()).thenReturn(null); + Mockito.when(backend.getHost()).thenReturn("be-host"); + Mockito.when(backend.getHttpPort()).thenReturn(8040); + + TNetworkAddress addr = invokeSelectEndpointByRedirectPolicy(loadAction, request, backend); + Assertions.assertEquals("192.168.1.1", addr.getHostname()); + Assertions.assertEquals(8040, addr.getPort()); + } + + @Test + public void testSelectEndpointByRedirectPolicyParsesIpv6HostAndEndpoint() throws Exception { + LoadAction loadAction = new LoadAction(); + HttpServletRequest request = Mockito.mock(HttpServletRequest.class); + Mockito.when(request.getHeader(LoadAction.HEADER_REDIRECT_POLICY)) + .thenReturn(LoadAction.REDIRECT_POLICY_PRIVATE); + Mockito.when(request.getHeader("host")).thenReturn("[2001:db8::1]:8030"); + + Backend backend = Mockito.mock(Backend.class); + Mockito.when(backend.getPrivateEndpoint()).thenReturn("[fd00::1]:8040"); + Mockito.when(backend.getPublicEndpoint()).thenReturn(null); + Mockito.when(backend.getHost()).thenReturn("be-host"); + Mockito.when(backend.getHttpPort()).thenReturn(8040); + + TNetworkAddress addr = invokeSelectEndpointByRedirectPolicy(loadAction, request, backend); + Assertions.assertEquals("fd00::1", addr.getHostname()); + Assertions.assertEquals(8040, addr.getPort()); + } + private Object invokeCreateRedirectResponse(LoadAction loadAction, HttpServletRequest request, HttpServletResponse response, TNetworkAddress redirectAddr, boolean isStreamLoad, String dbName, String tableName, String label) throws Exception { @@ -398,6 +453,22 @@ private RedirectView invokeRedirectToStreamLoadForward(LoadAction loadAction, Ht return (RedirectView) method.invoke(loadAction, request, addr, forwardTarget); } + private TNetworkAddress invokeSelectEndpointByRedirectPolicy(LoadAction loadAction, HttpServletRequest request, + Backend backend) throws Exception { + Method method = LoadAction.class.getDeclaredMethod("selectEndpointByRedirectPolicy", + HttpServletRequest.class, Backend.class); + method.setAccessible(true); + return (TNetworkAddress) method.invoke(loadAction, request, backend); + } + + @SuppressWarnings("unchecked") + private org.apache.doris.common.Pair invokeSplitHostAndPort(LoadAction loadAction, String hostPort) + throws Exception { + Method method = LoadAction.class.getDeclaredMethod("splitHostAndPort", String.class); + method.setAccessible(true); + return (org.apache.doris.common.Pair) method.invoke(loadAction, hostPort); + } + private HttpServletRequest mockStreamLoadRequest() throws Exception { return mockStreamLoadRequest(-1, null, true); }