From cfb8cae7224b3dd8e00528832c15d5b8ebb75486 Mon Sep 17 00:00:00 2001 From: gaoran10 Date: Tue, 15 Jun 2021 01:19:42 +0800 Subject: [PATCH] replace `SslContextFactory` with `SslContextFactory.Server` --- .../handlers/kop/KafkaChannelInitializer.java | 2 +- .../TransactionMarkerChannelInitializer.java | 2 +- .../handlers/kop/utils/ssl/SSLUtils.java | 11 ++-- .../handlers/kop/KafkaSSLChannelTest.java | 60 +++++++++++++----- .../ssl/certificate2/client-ssl.properties | 18 ++++++ .../ssl/certificate2/client.truststore.jks | Bin 0 -> 869 bytes .../ssl/certificate2/server.keystore.jks | Bin 0 -> 3821 bytes .../ssl/certificate2/server.truststore.jks | Bin 0 -> 869 bytes 8 files changed, 70 insertions(+), 23 deletions(-) create mode 100644 tests/src/test/resources/ssl/certificate2/client-ssl.properties create mode 100644 tests/src/test/resources/ssl/certificate2/client.truststore.jks create mode 100644 tests/src/test/resources/ssl/certificate2/server.keystore.jks create mode 100644 tests/src/test/resources/ssl/certificate2/server.truststore.jks diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java index 587c3c5db8..4f2c10727f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java @@ -49,7 +49,7 @@ public class KafkaChannelInitializer extends ChannelInitializer { @Getter private final EndPoint advertisedEndPoint; @Getter - private final SslContextFactory sslContextFactory; + private final SslContextFactory.Server sslContextFactory; @Getter private final StatsLogger statsLogger; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelInitializer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelInitializer.java index 8a07e93adc..6cf6912fdb 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelInitializer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelInitializer.java @@ -32,7 +32,7 @@ public class TransactionMarkerChannelInitializer extends ChannelInitializer sslConfigValues = ImmutableMap.builder(); CONFIG_NAME_MAP.forEach((key, value) -> { @@ -114,8 +115,8 @@ public static SslContextFactory createSslContextFactory(KafkaServiceConfiguratio return createSslContextFactory(sslConfigValues.build()); } - public static SslContextFactory createSslContextFactory(Map sslConfigValues) { - SslContextFactory ssl = new SslContextFactory(); + public static SslContextFactory.Server createSslContextFactory(Map sslConfigValues) { + SslContextFactory.Server ssl = new SslContextFactory.Server(); configureSslContextFactoryKeyStore(ssl, sslConfigValues); configureSslContextFactoryTrustStore(ssl, sslConfigValues); @@ -226,7 +227,7 @@ protected static void configureSslContextFactoryAlgorithms(SslContextFactory ssl /** * Configures Authentication related settings in SslContextFactory. */ - protected static void configureSslContextFactoryAuthentication(SslContextFactory ssl, + protected static void configureSslContextFactoryAuthentication(SslContextFactory.Server ssl, Map sslConfigValues) { String sslClientAuth = (String) getOrDefault( sslConfigValues, @@ -248,7 +249,7 @@ protected static void configureSslContextFactoryAuthentication(SslContextFactory /** * Create SSL engine used in KafkaChannelInitializer. */ - public static SSLEngine createSslEngine(SslContextFactory sslContextFactory) throws Exception { + public static SSLEngine createSslEngine(SslContextFactory.Server sslContextFactory) throws Exception { sslContextFactory.start(); SSLEngine engine = sslContextFactory.newSSLEngine(); engine.setUseClientMode(false); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelTest.java index ff2685bb82..350a1a3f43 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelTest.java @@ -36,10 +36,13 @@ */ @Slf4j public class KafkaSSLChannelTest extends KopProtocolHandlerTestBase { - protected final String kopSslKeystoreLocation = "./src/test/resources/ssl/certificate/broker.keystore.jks"; - protected final String kopSslKeystorePassword = "broker"; - protected final String kopSslTruststoreLocation = "./src/test/resources/ssl/certificate/client.truststore.jks"; - protected final String kopSslTruststorePassword = "client"; + + private String kopSslKeystoreLocation; + private String kopSslKeystorePassword; + private String kopSslTruststoreLocation; + private String kopSslTruststorePassword; + private String kopClientTruststoreLocation; + private String kopClientTruststorePassword; static { final HostnameVerifier defaultHostnameVerifier = javax.net.ssl.HttpsURLConnection.getDefaultHostnameVerifier(); @@ -56,24 +59,48 @@ public boolean verify(String hostname, javax.net.ssl.SSLSession sslSession) { javax.net.ssl.HttpsURLConnection.setDefaultHostnameVerifier(localhostAcceptedHostnameVerifier); } - public KafkaSSLChannelTest(final String entryFormat) { + public KafkaSSLChannelTest(final String entryFormat, boolean withCertHost) { super(entryFormat); + setSslConfigurations(withCertHost); + } + + /** + * Set ssl configurations. + * @param withCertHost the keystore with certHost or not. + */ + private void setSslConfigurations(boolean withCertHost) { + String path = "./src/test/resources/ssl/certificate" + (withCertHost ? "2" : "") + "/"; + if (!withCertHost) { + this.kopSslKeystoreLocation = path + "broker.keystore.jks"; + this.kopSslKeystorePassword = "broker"; + this.kopSslTruststoreLocation = path + "broker.truststore.jks"; + this.kopSslTruststorePassword = "broker"; + } else { + this.kopSslKeystoreLocation = path + "server.keystore.jks"; + this.kopSslKeystorePassword = "server"; + this.kopSslTruststoreLocation = path + "server.truststore.jks"; + this.kopSslTruststorePassword = "server"; + } + kopClientTruststoreLocation = path + "client.truststore.jks"; + kopClientTruststorePassword = "client"; } @Factory public static Object[] instances() { return new Object[] { - new KafkaSSLChannelTest("pulsar"), - new KafkaSSLChannelTest("kafka") + new KafkaSSLChannelTest("pulsar", false), + new KafkaSSLChannelTest("pulsar", true), + new KafkaSSLChannelTest("kafka", false), + new KafkaSSLChannelTest("kafka", true) }; } protected void sslSetUpForBroker() throws Exception { - ((KafkaServiceConfiguration) conf).setKopSslKeystoreType("JKS"); - ((KafkaServiceConfiguration) conf).setKopSslKeystoreLocation(kopSslKeystoreLocation); - ((KafkaServiceConfiguration) conf).setKopSslKeystorePassword(kopSslKeystorePassword); - ((KafkaServiceConfiguration) conf).setKopSslTruststoreLocation(kopSslTruststoreLocation); - ((KafkaServiceConfiguration) conf).setKopSslTruststorePassword(kopSslTruststorePassword); + conf.setKopSslKeystoreType("JKS"); + conf.setKopSslKeystoreLocation(kopSslKeystoreLocation); + conf.setKopSslKeystorePassword(kopSslKeystorePassword); + conf.setKopSslTruststoreLocation(kopSslTruststoreLocation); + conf.setKopSslTruststorePassword(kopSslTruststorePassword); } @BeforeMethod @@ -109,7 +136,8 @@ public void testKafkaProduceSSL() throws Exception { String messageStrPrefix = "Message_Kop_KafkaProduceKafkaConsume_" + partitionNumber + "_"; @Cleanup - SslProducer kProducer = new SslProducer(topicName, getKafkaBrokerPortTls()); + SslProducer kProducer = new SslProducer(topicName, getKafkaBrokerPortTls(), + kopClientTruststoreLocation, kopClientTruststorePassword); for (int i = 0; i < totalMsgs; i++) { String messageStr = messageStrPrefix + i; @@ -143,7 +171,7 @@ public static class SslProducer implements Closeable { private final KafkaProducer producer; private final String topic; - public SslProducer(String topic, int port) { + public SslProducer(String topic, int port, String truststoreLocation, String truststorePassword) { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost" + ":" + port); props.put(ProducerConfig.CLIENT_ID_CONFIG, "DemoKafkaOnPulsarProducerSSL"); @@ -152,8 +180,8 @@ public SslProducer(String topic, int port) { // SSL client config props.put("security.protocol", "SSL"); - props.put("ssl.truststore.location", "./src/test/resources/ssl/certificate/broker.truststore.jks"); - props.put("ssl.truststore.password", "broker"); + props.put("ssl.truststore.location", truststoreLocation); + props.put("ssl.truststore.password", truststorePassword); // default is https, here need to set empty. props.put("ssl.endpoint.identification.algorithm", ""); diff --git a/tests/src/test/resources/ssl/certificate2/client-ssl.properties b/tests/src/test/resources/ssl/certificate2/client-ssl.properties new file mode 100644 index 0000000000..467115d16e --- /dev/null +++ b/tests/src/test/resources/ssl/certificate2/client-ssl.properties @@ -0,0 +1,18 @@ +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +security.protocol=SSL +ssl.truststore.location=./src/test/resources/ssl/certificate2/client.truststore.jks +ssl.truststore.password=client +ssl.endpoint.identification.algorithm= diff --git a/tests/src/test/resources/ssl/certificate2/client.truststore.jks b/tests/src/test/resources/ssl/certificate2/client.truststore.jks new file mode 100644 index 0000000000000000000000000000000000000000..dd152ebfc21d412909ebeea0f4cc70053b8eb7d5 GIT binary patch literal 869 zcmezO_TO6u1_mY|W(3o0$%#ez`6WPZ6?fnoUk26)JyQcq1_ov|gC=GrgC-^}CQgQJ zMjtEc|NAx=@Un4gwRyCC=VfH%W@Ru4Hsm(oWMd9xVH0L@_JeUac$kv1U>tU24i^to zu%9810T)OQHxF}iey$-u&`=_nbacN^| z6QdHc=NMTTn41{+84Q{jxtN+585yo9FLUCZ)ueDCQrXvMwP}=J$A_}>_46BFHk|r3 zUy=Wfz0Zr*g0SB&{TYuumeuaNYWKb2ltPaD>Fz|HCG|f|?VfypB+%CIKkS23y>oQ- zzuj??vu5m$`jfY3@7>*eS7fek{QS?S>+80)2AlSGgv>c3BsW5eGx}lzn|XWyzys>n+%6sCr=kgHWM==10!+>0iz!nLW~SMly$pq ztTle~Ekk-{!N%jR%vXCZDRer~zWn;VX=nSH1_ zq+KyR+(~`w5w5dx$K;MK$l19@aJ^XU@(XE-OUu1m5`qp)a9kfa{pd}@b*rzI-VKO| zn%2oMVSbt4&yq==8I@liN1MHE3}2r9Vs`hxg$usdestb^-mJy%@6YEqw|aS&+`S`M zCG)I&Kl4)`4jnH!C&t=K6>FN8&*PaM*P`65bv}Ld^YWLOGGgpg7m7Gtm&iP+A3Zm; ztyl5w2J;|pCsV_c;H5RZrwm(qvMePwTvhCrHudJ`|KnHc-yZya){f@(&DY)@k6jpB eoFJ9Ma6C4Hv%&QgyQphWz(?;Tr%k6dXafMPc}l7P literal 0 HcmV?d00001 diff --git a/tests/src/test/resources/ssl/certificate2/server.keystore.jks b/tests/src/test/resources/ssl/certificate2/server.keystore.jks new file mode 100644 index 0000000000000000000000000000000000000000..7517baa590672d7001799269ba52f319e4e25501 GIT binary patch literal 3821 zcmeI#XH-*50tfIkT0#?~Nbev;5~>9RB2Pf7^r{pkfFP34YX}$!9R#E|X-WsB2ukmQ zh;%7WP?|)F6lrnu-dW#q&z`eqzwCZ`ALhgV-kCdR=FZ&v8xoO31ONb_zY3*;9opU9 z6954E!wj{>02JmILUX(XJpC_;5;B+A6 zcDq9LQv6ec><`YZ3p&}6qX#pp_h7&#w6K-eFq1_nsYQL^N!8MmZ&AGo!%uKAq*DC&D4Fc1Lzb%w~V zpFBfAK%<~k)KrD+#=aA0`h%L@-@xO&?}ncBy~>~K4$2LgSTJP`(m?s*PL`nzFdo8m zb#HlQM}TnzKgtX#`l1Uu#M8smo#)zAex{1mrewsCzsN^B+}5xyL9OcM#cX_nWW|9x4UGeKnyi2L9J#+lk)Qpg0`>Ky;TU;J3U;PesM#~b8hyG zKi9HPEBLbxwd4&RBrs^yw>+#QhxYNUaKUKtx2UoepEYMLR;rZzv&czy=YDCcY!iGe zKe1ZT5QdaTcp4Q2(hVTOV_Yt=SC8{YbINNoFr4Ul=|>oCW;BLH)J+h3ZSrj}wj8eH zeFOjkXk6VL>|CAPG5_30_YDC63J@FzA;IY>p%Rop2v8h!21s625FBLv`nlu7oLQ*( zMSVTfTAD9g7}m$%1=<_#&=N=Aac1t=Uj>@ z^t10V6aqGVah6Tw2zjM=cCYVCEY!Bcat4*^0nSKrKa|~X8YN%TS=gmQcKSuA}6F!STaTF zW1XV-X%%P16QWtOPkX9nM>u3^L;DQ3*kfhq6+x~`LwvGwMFdJlK1Rk)u~OB-HLx|w z`9b*&+5n5a*^p%WQSJUA5PDd%I6w@Wi@Yvlo9h1~M@X3Ip0L(-P+3=|M#by3Y6V5t z$;92p(*y^+MZ^mR+Pyp(et zq5Q!Z4g2PrOoF>m62#SLOdGs;i+C4!TY0&g4>m>5_ak&-Q9DOzJ#t?T-_EX>HGMlX zbdp&c`-w_O;ZDKYYWFe|@(X>3B9&haY+YprX}IH8@`#?#$ubU%sYL6hNy0D&`+|`d zZtN~(>{>drP;I^3H5YA|a;_e*zTR}zpm-RSviRs*D0Fpvb6#X}$BqUJ4rA#FR4BB_ zHktRsi`~2qyn35Oe<44x+GQzRS(#lCYBXOiTV^j{LL9+Ps*)gEJezhq^*h8B16TZE z-t!jr27XAJRz~z(C}_ft^E!IZsi-p%VDKR#vD)SX!jbbcYs%W5q9k#5bHyL=22E|dehFzPwt$I`ZSir5aC24fsDdIu0QO?5K;5Te+f^0S5-aXA zVLkT1tSDH*S&93G32sj_?dR1uIwxi*XeCHg=TwUOgW?_LtUQn~C8u=#q-oWosk=p5%k4HQev@#ib zUxCg6>=@D2R=lxgRqmxqMM1>oraPvegr@s8gc38~tiXruB0FxX%$=}Xq&Qf*tD9rq z`*ksfUYMZ^)Aga}k#7CGM#oin_jX$yRm23>qmF)XYI1+^!W`Q&2erST@OM3vI4L?k z+u>1u-_zhH+`*Gqo?UJ(6va?sns_(lvJ)h3d5PxkjP{F82^sk4G7gT(0+|BV3b8L^ z+x;qOYl~}ZHl3CH*{H=4s+8i7dia9nY(~w+nVwZ?Cp_;6)-eq8J%{EMm)&AF>$TAD zn#vQr{?h1AW&9-y39=|e$)aE&i=tI3eStzIfbG{`3dR3J6qjV+zl!3I2si(U6*6PU zK4Bz(w8i)EN8F?Z44chO`%{AX5d=DPS? zn2ERGMM^zTh_ZMUt_C!zw+2AL(WbhVic5iKxXqC#c~z*wcq{yj&Am+72d?3z8f( z_Ea8^uVC?Ztn^bt_4@cajgPZ} zz8&1&*CJIzu-O)i>enXfZT5De~K@+o*K@$@f6DPwq zqmLE!|9u+_c-c6$+C196^D;7WvoaV28*&?PvN4CUun99c`@uLIJWNShFb+F1hl__P z*w2v1fD5FDn}<0$Ki7~SXef}!%)?fkrw0@;kQ3)MGBhwVG%++YG&eGj0&7qANIki-Z?t^ z-|jfcSu=J={mI+2_wH`KD>7F%e*Wjv^>y1?gH8K8Lgt(ia%lT@VcAc~X_tQ8d+NaA z-1drl?FFUjiu%txSbGa=+Zdbg{u0fXPcLg*D7WCy_rIs6#ve4>@VI|vz}(xJ`|tbA z?YQ*)mP0l-3Rlc$R#n~9l`fe|@`fYA>OAx4HB%DP=Q z)*8S0mLWZ}VB>LD=Bvjqo?Uff@4~xBYvz-cG^ye8DKet@x%?-!b%sx~d z(yo{u?xeo;2-jJ;V{%6qFf zP3vTsFu%<2XUQbbjLI*Mqs`tnhA&TlF}wTU!Uf-JKRRzdZ`R`X_viDQTfIC>?%ols zl6h9XpZTc|hmMz=6JzbAiZ#v4=kZLBYf