From 722ab9710fe04fb50ce68d06bf19ebc560305a13 Mon Sep 17 00:00:00 2001 From: Administrator Date: Wed, 26 Jul 2023 21:30:00 +0800 Subject: [PATCH] =?UTF-8?q?=E7=A8=B3=E5=AE=9A=E6=80=A7=E5=8D=87=E7=BA=A7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .classpath | 4 +- .settings/org.eclipse.jdt.core.prefs | 6 +- KLALB协议规范V2.0.docx | Bin 19158 -> 19427 bytes klalbclient.ini | 3 + kserver - 副本.ini | 3 + kserver.ini | 2 +- linetable - 副本.txt | 19 ++ linetable.txt | 20 ++ .../cloud/network/klalb/ByteArrayRecycle.java | 28 ++ .../network/klalb/CompressedSocketBridge.java | 34 ++ .../kne/cloud/network/klalb/DATATPacket.java | 21 +- .../cloud/network/klalb/KLALBController.java | 10 +- .../kne/cloud/network/klalb/KLALBPacket.java | 7 + .../network/klalb/KLALBRemoteSocket.java | 25 +- .../network/klalb/KLALBVirtualSocketImpl.java | 301 ++++++++++++++---- .../network/klalb/LineDecitionComparator.java | 24 +- .../kne/cloud/network/klalb/LineManager.java | 28 +- src/org/kne/cloud/network/klalb/Monitor.java | 129 +++++++- .../cloud/network/klalb/MonitoredSocket.java | 19 +- src/org/kne/cloud/network/klalb/SendTask.java | 7 +- .../network/klalb/SimpleKLALBClient.java | 15 +- .../network/klalb/SimpleKLALBServer.java | 9 +- .../kne/cloud/network/mport/SocketBridge.java | 40 ++- .../kne/cloud/network/mport/StreamBridge.java | 14 + .../kne/cloud/network/mport/ThreadTool.java | 2 +- .../network/mport/VirtualServerSocket.java | 6 +- .../cloud/network/mport/VirtualSocket.java | 15 +- .../network/mport/VirtualSocketImpl.java | 19 ++ .../cloud/network/nathole/NatholeTestC.java | 32 ++ .../cloud/network/nathole/NatholeTestS.java | 36 +++ 30 files changed, 720 insertions(+), 158 deletions(-) create mode 100644 klalbclient.ini create mode 100644 kserver - 副本.ini create mode 100644 linetable - 副本.txt create mode 100644 linetable.txt create mode 100644 src/org/kne/cloud/network/klalb/ByteArrayRecycle.java create mode 100644 src/org/kne/cloud/network/klalb/CompressedSocketBridge.java create mode 100644 src/org/kne/cloud/network/mport/VirtualSocketImpl.java create mode 100644 src/org/kne/cloud/network/nathole/NatholeTestC.java create mode 100644 src/org/kne/cloud/network/nathole/NatholeTestS.java diff --git a/.classpath b/.classpath index 7a21a4e..0039a3b 100644 --- a/.classpath +++ b/.classpath @@ -1,9 +1,9 @@ - + - + diff --git a/.settings/org.eclipse.jdt.core.prefs b/.settings/org.eclipse.jdt.core.prefs index c59d0c6..85dd98e 100644 --- a/.settings/org.eclipse.jdt.core.prefs +++ b/.settings/org.eclipse.jdt.core.prefs @@ -1,9 +1,9 @@ eclipse.preferences.version=1 org.eclipse.jdt.core.compiler.codegen.inlineJsrBytecode=enabled org.eclipse.jdt.core.compiler.codegen.methodParameters=do not generate -org.eclipse.jdt.core.compiler.codegen.targetPlatform=1.8 +org.eclipse.jdt.core.compiler.codegen.targetPlatform=19 org.eclipse.jdt.core.compiler.codegen.unusedLocal=preserve -org.eclipse.jdt.core.compiler.compliance=1.8 +org.eclipse.jdt.core.compiler.compliance=19 org.eclipse.jdt.core.compiler.debug.lineNumber=generate org.eclipse.jdt.core.compiler.debug.localVariable=generate org.eclipse.jdt.core.compiler.debug.sourceFile=generate @@ -12,4 +12,4 @@ org.eclipse.jdt.core.compiler.problem.enablePreviewFeatures=disabled org.eclipse.jdt.core.compiler.problem.enumIdentifier=error org.eclipse.jdt.core.compiler.problem.reportPreviewFeatures=warning org.eclipse.jdt.core.compiler.release=disabled -org.eclipse.jdt.core.compiler.source=1.8 +org.eclipse.jdt.core.compiler.source=19 diff --git a/KLALB协议规范V2.0.docx b/KLALB协议规范V2.0.docx index 0375c536f051b23c83f6854a79730b1bd6009736..0886dc081bf071ea9b7964be24c03d05f788f572 100644 GIT binary patch delta 6334 zcmZA5Wl$7=vjFg;K{}O1e7)p6|Vx z_n&!Rc4u~WclPW4=F^e3GLce*u#8~fVRIw^!0WU0GbHDsE=aC9Z!RUsEYP+IE&N0| z_t_5gL%WLX8T7J1U4I-|)tqyVx2Gv~2G>1&@3|9<@a_JMPA8L3oyvUkX9$bd4%#mc zrYY6p1*d5iTslgQj5uc7Iz_-FskYt>I(0ESW+?_&n65i8t<*>jYCn|R%uRb}tD~Tj z0000?2q0Q>Q+9~coLv6jMJ~ty0O7x~t%r@aw}+<>uZ@Sd9Z!IpYhJ;ro z0Unb3l_oc$xgVw7o8DNb0~_ALy~xZG>;HJ$Y;F`y+<|j?DFyZ!&I~eFv0U$HD3cXr zU`J-E*0(n;%5)}06b)mzGuPF8M5gH}v4em^(&MAB4%gwo(ap5iOD1e3iQj$nXAs1( z;q@ zd`PhPE*f{e@x{kVbwahKuXuk#a9#?e?uEk^cshsL?3Ua3Tf5{?>wNj)94zWwt>FwY zLe|b!=W`v3K%tdY&JUw$kzlF&{iM(&Z>KiR#d?`lM#_+&t2`W(q%|vwz)t2l95Jm( z*O3;btv>wqX1}m>N4mQG^$L3mC>n9b0yLQ0n(8)x&6jUr+si*<9}lZ#lgr+UKh+6E zfia-;Vq-1*_)|QF1^_fe5Ai4|zakdvl9eqt1|8pj>>w)@FFHGp)_bpI`OZ<-nW zAfo=Q@VkVXeZxwN*3Kyp=8! zVd&xoj#hCBz8IEHf5~3CItvo zFCwYkMORkIv`zcLbNq(X=*O#v%S${ARzbKD7U~8v3bxg7%?2@zJFko^gXF)krBkn< zk{jAFGcG2u071!+Z2Mv6k`uqDVDjV275nC&Z`DSZdpb8E55w0nfeax_0eg z(yZTr|E%G_!?AJs{h5n}%+Xb1Sb_+&py2gAx5P*%ZNLzu-AQWAouJ5y+jm|y-@RJG z#0#USv(j?Ozw4??LRs0$mGwegPP6#_it+oBtk16ZK5En zpXkJTBHxayw6>on)9{LZ;_9n;c_F8AHWc&&8TGGJQ0esg{x!0R@Au})+?Tj(GmtGf*e|?XmEyvc(VB;?;=*z4$RdiU9YpKty=cMTvZ79@ z5d0dBX@AJFCHuM#hseGT$M@iKTQEh~q-w|GFvdL4HZY`-ktqFYqam4Ix!#{$tLs$O zA!0})ADS>}Ibq9E$c&0j$E(<0+Nq~j)L7tD~uGrG_PCfUsj%m|RA zL}yRtwg_UFP(3kmJ{TCj_97`5+8U_PLXkA@Mg*kjZbtsp)vJrKo-83l0KKrh%+akSIr3tCMo7uu)}TZiag!h_Rzm(D*GH29Y|`+GSQ0O}PQ zM-u}zocXq>@Z#wpiPg8c6dfJZAG_?{>Gh-vp(uP`F8zJT-$gSk4APLwyV9f{QHvah zTs+0|z|efIFKj)YW(r0I{$w_pFp3`ftt`}^%AE4d;44!~FC@OtRblG3F&A-f$;D{&ri4{9 zFqLf>p(PJ(4DIyh0zG=xGkc@xss(!uk#Bv0I9r6RW&N>H4wTZr)|~_|J`;}@Iz7iK zJ5Wkm;E&R9saIAEFu6ut{+eoOaN7Q8v>@R1C-3thGfw+dB71JgK(X99?k_}4Q&c$& zezxEHbVHXy3AIu~g-wwtqQ5fjY^G1d2LSZjp|n&KkYzn$6G{l8@`+9tY(&I^SB73r zz^jhCrTYceWMJAf=Sw;5Qr=$G&ag>5hJzfpo&Z86MRRc0_oPp2_J&@8{)?yg6G zG(l_(SSeYG^t(vYnLdY?tD(CRnWv+5JVWU{^|o0JW&WuBDhj<}U*q2E8)8U=zvj&o zo2JEqWMvd7qbcc?CauTVqIZj*-*$&&$0!PUvGY_I0^+pCRpnnxg!7F~GrfCPMh@K! z2MhAVN7*~0%b32hhp`=9&WypvlxF^cH@Yx0 z#MI7IZW(m6Y(Ih!GdQ5vEo7v3z;#o-%`F)%G<@F@26!Ye2jv z_Jw;Sx%=a0XOPLT7>8TBGkzA$nHn$oa=?*Fy^q}@*jwdVW<}4V`~$;aoZx}6iA$NCwTTGYQG!@hHN zPGGwa;x!HgDs(E#CiV4O$&zl$x?h}4Ud<{3PD)bW$amFE=ARUse^vSE)9-*Qmki&e zKq0ch$Nt=#$0n5%0QXJQ&UloB#5cD;*f^`r{jiqxteBbIaH)4oYm``csc6~>eif%_ zHOxLs*bWKfiJwnl_?XvYfcm^*-`3PxTW z2l5TQVC@*M`eFI`TA<;JRq}>kPmJn$C|1_+UQ#n-O+<_nMJmhSnKvL6Qh*`Oi0fl9 zKsFiM_ztV9MP**rH6Sl;m%j*vSH6r#vz8406yD|V^JAar zk~NGX)soN^XJ$OBl{k>{K_j7Sz%lVDQbnmP;~gb$B>C052J`j+~^c^_o>xnP6-u_#6U0LTv(e z*6NBjGX9`RQ9auYQSP?sQym=2OK;NmUJ7kfEgQ;r*wV%#Ne6pEHM-Pm>dJ*r>&5k` zwDM#;N}I_9Y>c}f?Laaqi6-oBMxhrhFjt)n0}>w+o66YVEQ<~(+MJF%sCyvodZ{p% zX2fGGfy*}i&sBZ-H)#8@=0b{ImiwrK*f;qLH;JT=aP$<2uTkKMxy<&hy-cYRB=ao< zw*Dw^>n4>GInrMMfgjPuO-)mS{U+OP@+HYZxhrhVK5?pya{JET6%Doy~y@K-Xfeut{K<#k@1#{Gq zucvTBqtqa1&WIRoH%j@p@?K0Pfw}k3nFl&yo5NhyBri0y*~=8lGwUCs$W!}iWCfq zO6N5!XT+~OqwsBRwgj@|t=pCrnJHkLXXG=O_W@bD2cTd&Vt8>j5EOYL&v1!pG=QF^KN>Lb`jc@R%>yuEbK4NlIu%@{RhMC)SgpAKVSBY&C9!zh2u8t_xvw4LUt028g)FKRbSmeYXJ|>(;3fY~ z5qfGzeyAVSJfu&<(rI61by5m%UTjS08CoM=E2j22J z>=d!fYjigTh zzV)pW5IVi^`&tNt2)AW67AZFPZ%-#7p$DA90GV2;UI@};vOD@(cmDK@YEnVB{2IT8 zT1k*RP82c3r&#gbC?f*Y3$^v5-lEgpySqFaqm=r z#RPUfwTvKT=wcN0B6piOo#3vxWvl-_yN+-g);i(k#Y!vh%o2D>k2zxiLxp3*dDgu1w zWm75G?dgGbYxpntxOyk(5KLjN5V9?)yWpSH6#4JfgE_m7CH)TmG86@d;)HkR#|CRR zJEjL2VyW@ppTB=Nhwz&`0XER{VlCm12mo*!o$k#$pCn>Y&2`a_TM`rHk+NSY#9I1j$J4+<1`Zx)wa7uQ@gnF!vP|+A^W=MSB&)mymPSvDY$z- z6<&Z6=B_VhZDX;DLtEA_lr;U92-iUBtaeX3N~QsOTCUPmBNGDpM|cG6Vks2nIH_e< zz>`wiE{dh}oq1H(Hdy;Z>u&f=vDkg*mZE>ok8crw1N`w*3f()SEf*y*2#c&#Q!LcP zCx?zqMb2;nAu%Y|1js`_HsXI{@%j}DE81K0QG%!HjMqd*Xb-7k(Ho;L&i2QnZkWzDk`Lj4M?Jtb3!lj1oY10gG!eTKSv`%lr~QtKUalzt zJE~t>5SP)j6Q!Av5y5M(ml6JjxjEp+$IPxmEu)1Q2sPc_FA&R3Yo~t7+?aJLv%?d8 zRS{edP7!2Jp^TW=6Q%kD%{|*@QyuiXms|KVu5=gXE`%1A(FRjWoixzg{-kN!lrntN zWcPTY2sa;?qcOuBj&iz|t=6zCebfq;)!d>!l?r$0f8%=QZ6oTsgJOrVIxV|MUtSAM zh~xeSvG+{Bf5?#xP6JEFIynwrZN|YG^e@bw4s>Z&nWcN%CLoH-5Jg1V>Cz2t(f{ilAAKAGuQheC`r;3e5U&lor!`m5v3l-N6z|jpVH^Bwp zp7ga=rP?lrnYWM69Bub`9T1y3UimFtHKgWjz z^yEgrl?Yb?HLlB}MR-0Ni$m0?_SD>;Rrl;A$em!=wrIR2JOP*4VFSWe@rN_KwNR>V z2Ud64wHyId6_D>wU5U{Y5u3iw>XlK-4^>7$t6Hjd9h$h?X@Fw`e{}$6WUZ0v5fNlV zygjOSVb_nmSco!g5(6Yi#4|HdG=E7VSb$z>9)~?xDyu7849_VY`Q_8 zsh(KEj_1M#Cg*yG+{75e=`8xZ{d)e%>k0nXyprR zQ5H8pnXM^wMB_~{`}kc$rq&iFg3*R;lmgP8h&<%!zCnN*38V zUQbc-Fqf(L_GXVoz)$Y$D0B&kYlX|6*ThyC>wrrkWt{Wf1xC}AgW*rY5Gnuw^6<8WPASp+qXazUc6q3|bsmX(xRIz(7t>2{Yck#3Q0SVFo}knToF0pa=P zy?Otc_vL=O_nbLj?wsFN8v1++dXgxfAasopjE;up`AO~(CY6!!u-CzAKJGXuD=1sL z1OWLKj6Ky#zb$uF+OQlO#1&~VWrbv6}Lua!Y@JIO6YO1p$ZjaBtf)OD_<}@NF$&yUpN#XU;_8E_K0(~ZbE5%KFLB$s{2E(k6 zAs5s8bUU;lJO|gX**3b4>a+s}HaS5)6DQ#8J}!)rFLOC7#(X#)Zg0=q%#7lgDwN%6 zRJkYN?O(i*<=m}P+oy4jH{k6nT850=;w+={<7PmxLhz_)PQ%M`;v5X|<+x5toNz+c zY0UKU-~;}tTTEiX29FaVEG^6_&s6Crfb$(*1VQ*5ck4TtK5x^!cl2IfPyD(+Jg}8; z5T@h@(oIw_2D0z5-d5bswkG&CLL2@x%fdf9bZ>@u+Xk!3RAoqi$b zP^tNTYoXdwzD|?9S-v$N8JC`y#om!`sb==Lv%Xh4SS#l~Clb0E=#mx7QGI6?lLD`{ z*b9w)ZPegrA_+c?9hq#|7bT~3DqQE;8w}W6Szdz$c$!>WI58gyD1_~&xwLrrnE8(M z=L`6@+x~^695BoF z5u)R#V{dKpvdi1w2hd$T`2yeVQJWTM%Ow1J_$J{ZvvOBOTTtr!V>!^xU7W!G?i0DL zi^$CLr2Ch{Lo0B@If8+SbbA^pHqF(Wq`XGj6&F(&pJgAj2m3zTtTD-9e`x?j30>MF za%Sjvr{q?XO~B};ONYKz{QC%ipRMvj$-U)>n002ypfm%Kab$HDcfU6ks&a%^wyR(Y70BBZUlP7rt1)vrPBI5Y)z|$7=&_7Tz;#6n5v5$>z#5#UpjF!hNPXEZMl&9_(d}BpXn!6j`XgFA zl1P4>p>F?E5Y@e`BQ)9aH2oc$gjAU>C>#MB2TzP#&cd!|F>xK^!awFKl@G6E6=-%>PuJVaxE2C$$(8f9;W+7TmsKJlyzIZF#GZ^;S&Wlh0~Jh$ z4x2NP&nSRpTr(YE0 z6s`JU<9yv9i}l$CL&^lb(9NIU$=epSOZ&yLU$q!G4*|w|o?(R)RpSwo3Ec2eB~20p z3e@yvhPfXrrsUo$_VE5vvq7f9ZU;T1l{0k#@@gWnb3w;{7W=8*2bndKJg3)X#BmMb zoB9%Eb$Fjjqlf4@H8T-b4a_E?Gpp5ppN-(HpCuwxkGq(C9SsC9Xz;oh=;z)}uI-T) z(Wd7ViNStj`ja*;QYG=SfFQk?rYf7WWK+yxtL)**x+GYV-iq|P6w`UInnQIrQX3>eU|UwsE# z&)<*5wYM{n299JF2k+RcJd1Wpvs)X!^&wgQDRkwAkk;qqp%Cl55vU#Y8J=URarY@` zguh$y!y7*&iHjgLCRntBEzz!?i63QXc7OPbk9Bl3z+szCRdUXdN@r=?E9^d^bN`Y$ z4LrZ=u~)g^4(?@qY;V(&r*Umop>GgNlnfeZb^R$NPVBG)w2OWfq1|hYmDsvreX*`B z0qsVmp$M5Of#ZNHFG0*L($eByya;<5In?@)@boA0(a(@aJQaIM8X0clk>|WZZ-*E!h`8-!hkR##COsU_nHAXoMZv1+4BBVbi=0=wy(9^`C$5{Kc!!N| z+RYOEnFtMysvIdxPYW|QpjV{|S}Xrc%-uoHN`hGtYCyBhOSxGTms3@>R9-dZpg^8z zS8S&;c#DT{s}qtC>V(vwun{pXHqw@GtS~0>D~HSNu9}iIdlfPS+7JRz%B0|VuzcAP z7gzybV%m%Mhx4ZTc&6-Z-SoD@ zcboIWv&=sRF}1uF4X5^cet5LI4biz->Q^SvG%y53v|>~84by%=<`aslN}hbslEo;h z$YzO9QU|K!!JNTu2658|@$Wq*T_?KGl3r`RiCM{alEy-!offEl>fsCY|e9cXI5_P4&99LQc1Uq(Wqw< z46HL;0dMN#j|!l;#D}A6murUQfF%_Vop7nS-Df|OV9qHnST6E4P;*?^o?%(2T3|x( zrGti+cUD2nXJfOkA6NJYG<|6KE*hvmW?CAdW7W!&Y#5?k63TrCycDA?*Po5^i8CgO z_rK|gPIFO4ybOGtxARLG+rCozV|A8P>0qH23XQ$T$ZhT#gJbmt5P2D^;U%tvQyz{M z&L@1^VH}JMyM;Oe#Ae!)?PIrPDqzPniXW3KAf@DZ9jFP-_B#X!xi6j*2d_O6xcXetIcW#0HHkAeYacUm>;lt9Mf0VSlYyPrQnMwh~NIzpasQX zeD1@i+}z9VK$A)T!JDg{*ucWrS!S4jv+CAt)F618S?CnPFOhsvsJrp?sF8Ho-t&;p4syskGgYVt9!IT%%PCA2`MBa?w^Ld50 z*(3bd{>mqCrp!+PK#RkaLd;78q8LH6l4MozsChaOBzkTm>8Gm!=!K=}wkQ~9LDorNU2S2YCKbEP(qL447 zt=&~1!*FAcMh!{-ixsfm&6#n=YNxJzaxg}+mzx5mCKZ^KQhIk5;@{G z&q^dKtD2gfhaUU^8w{BOo8n-udD%_gfj^%Jsj5 zM4lu%&aJURucns*9LKARe5fk>$XFG(vdUHh z^wgm8+TlUD>RnoSE!1UI(gz%8EPR=_&s}));`erE+1q<4cauSye?4KB0o-v8A4iuK zHQ4M&jd>jGetr-lN>R|y81S8%jhwyqalDbM5^1P&{_8H#(1f7vy;=reufJT%fKL!#)v@hy$y^prV2|AxPJF8UdYJLk+gZ1wtl^wj% zJAxcm?pE8M81EQ6g2FnV=2*#WS1d7DFjHPNh^9}>8j)2!hAQXSfQB3qoKOYB#Ol)E z`=cpdEtdVY%ZjVn`L@e(_~mb!OKJDQsW!#_%|XY6Zi0Pmam6ASz}s1=RE&^?m8F@q zs&nGu-06x_!OeyrtMtpzH};U!kM}*tL!>fH=~gaiaizI;^>qpn=w;h7f!MWJo)y3mXjDq1rTX`DU*g?iopfC1Z1e-chv9^Cto$U@J1F`UPr=?_{ zsCY19OonReOSR}(+;-3>!dMU4LSoUamB<**ba-~&_&o@~{3x0t{l7=yzm(&*<2*gG z6gC{yCFz9JN!iE~i@Xvv*WflVEb)^Gmw zT47fJp=v@_j)l`dGbJd>njQ7$`)ke<>5Z-N=yx?E035UC%{aUw-#BuHkE?`Vg7>!*A zn=6yDFHXcul;qUB5cdt~m%hvU2JpfUVN*$CB^rXA)7l6F`wjrP@`ZPe7SaGFFuj0Z z$}%3X7NtpQt<0&0XBl%W=#HmMvSq`=qH?WKQb_b0e?IP*-_G_DD)pK%!p z*wwLi;+mh*I&$u)TbF(*oxvB;O2CL)ZB58S9ENEb z)rN3*QAc0G*T}b^=OP%>R*GYfFyYxrKi7VS*z)8t$smLSVW&?f&6oCVF&4 ztV_Pa3nqC%sL~`JCYwg^7hJk+ z@-xf9T_KACwe-+2>l*(!xNpZ->tooGdQAC4sA&pz9rq%+d%S)Js?_x?L-xkN&TI5L zWz)Rei5AK(e0LMNJ59HeKmuLJstX8ogmMKw{VK^hn11MKryMm$wNuW18xmAOp@Sq6 zx(e#Atyb~hMe1ZVZ7x?6P*-OI%J(?C?7Xc6M zXmpYAdd5W#xVCJ)5z!ii3=U}J`rv(yH=gU;?eo2a9MGTYm^SjkM0{&y?UCVm?{dqe zEb7?5l=fqi+3qh#Zb{J1AWS@sqbgo}w~f-%UFzZ^UX@(899s_^P?OaC>p8G5)Fq$x zPtF_)$NvR4i@aqF@iqs~!8-N3jg4T;r1MM1X)`8~vlW#2aE}n4?z)BOy_bBeRU7|o z#1_Fmsmifd+%egjqk7QlA91_i{LfFPdTC;J!R+k$2p5$idSa0xN#GE%hsh0> zwaiHbCHz9ay>t@YhB>m(bgHH9d2~Eut6h6@>A*$R0XBah}%1Y;tS@i*NYt)_nEAHaB@%kD~&S z41SmeeN*v-?>_%~EMCHR+ih0T2<~Q$bgKK-z`j*=IW7zQMKs1bri$m~``o2dq{*1M zLzW%B;>fW#`U5;U8<}~bOh{r;@F`b`2N0&XJajp-JiltdESx8r=W{n#F+0y`XeR{cs8ln1-gPfB^|sw^#uUT-Izx98z4~LGI>0I37LML8gnO z897i{F-q@TFTrepq9Mg=bJ>ZDNj0?a3#pb1i)k>d?ZbfX=#5BmT;3X`TXZ9@{-T>pVM_;jQKH1Jeh)}S7XEg3n$O3 zo+tAjPKsw^>k=x5?}ea+5ZT;IG@qP*GH*3^xfTt79AA#=4#%pRS)Ni!IAJuTz*x(b zT@H1kkD=X8#fAL)pk(Zw)gCgzNCi!L9gJKl7tZWPj>>b{%%}}K*&1UqThKuCxpJ4o z`Q(FTfeV?eDi7%4Mjr9dAE4xI%#uBL$Q zk0e$X22_KPit6Hk1rS1!@xKKa{?oYs2VVbY8Jivnq9;fGQ0E3z@FTG_d{tG@T Bkct2R diff --git a/klalbclient.ini b/klalbclient.ini new file mode 100644 index 0000000..b6ecda0 --- /dev/null +++ b/klalbclient.ini @@ -0,0 +1,3 @@ +#Tue Jul 25 19:03:09 CST 2023 +local=0.0.0.0\:35000 +server=127.0.0.1\:4569 diff --git a/kserver - 副本.ini b/kserver - 副本.ini new file mode 100644 index 0000000..944768a --- /dev/null +++ b/kserver - 副本.ini @@ -0,0 +1,3 @@ +virtualip=67c:72ce:765e:4db1:a02e:92fa:1959:29eb +bind=0.0.0.0:4569 +local=127.0.0.1:36555 diff --git a/kserver.ini b/kserver.ini index 7e44418..944768a 100644 --- a/kserver.ini +++ b/kserver.ini @@ -1,3 +1,3 @@ virtualip=67c:72ce:765e:4db1:a02e:92fa:1959:29eb bind=0.0.0.0:4569 -local=127.0.0.1:5212 +local=127.0.0.1:36555 diff --git a/linetable - 副本.txt b/linetable - 副本.txt new file mode 100644 index 0000000..320126c --- /dev/null +++ b/linetable - 副本.txt @@ -0,0 +1,19 @@ +43.248.189.107:65529 +cn-bj-bgp-3.openfrp.top:65529 +180.76.147.250:65529 +cn-ah-dx-1.natfrp.cloud:65529 +cn-nn-dx-1.natfrp.cloud:65529 +cn-wh-dx-1.natfrp.cloud:65529 +cn-zz-bgp-10.natfrp.cloud:23330 +cn-zz-bgp-7.natfrp.cloud:33336 +43.143.109.64:49965 +frp.104300.xyz:49965 +us.afrps.cn:49966 +hk.afrps.cn:49966 +la.afrps.cn:49966 +frp.freefrp.net:49965 +frp1.freefrp.net:49965 +frp2.freefrp.net:49965 +frp4.freefrp.net:49965 +192.168.0.233:4569 +192.168.1.233:4569 \ No newline at end of file diff --git a/linetable.txt b/linetable.txt new file mode 100644 index 0000000..cb2c257 --- /dev/null +++ b/linetable.txt @@ -0,0 +1,20 @@ +43.248.189.107:65529 +cn-bj-bgp-3.openfrp.top:65529 +180.76.147.250:65529 +cn-ah-dx-1.natfrp.cloud:65529 +cn-nn-dx-1.natfrp.cloud:65529 +cn-wh-dx-1.natfrp.cloud:65529 +cn-zz-bgp-10.natfrp.cloud:23330 +cn-zz-bgp-7.natfrp.cloud:33336 +43.143.109.64:49965 +frp.104300.xyz:49965 +us.afrps.cn:49966 +hk.afrps.cn:49966 +la.afrps.cn:49966 +frp.freefrp.net:49965 +frp1.freefrp.net:49965 +frp2.freefrp.net:49965 +frp4.freefrp.net:49965 +192.168.0.233:4569 +192.168.1.233:4569 +cn-he-plc-2.openfrp.top:4569 \ No newline at end of file diff --git a/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java b/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java new file mode 100644 index 0000000..c4490f6 --- /dev/null +++ b/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java @@ -0,0 +1,28 @@ +package org.kne.cloud.network.klalb; + +import java.util.concurrent.ConcurrentLinkedQueue; + +public class ByteArrayRecycle { + private ConcurrentLinkedQueuerec=new ConcurrentLinkedQueue<>(); + private int capacity; + private int length; + public ByteArrayRecycle(int capacity, int length) { + super(); + this.capacity = capacity; + this.length = length; + } + public synchronized void recycle(byte[]b) { + if(b.length!=length) + throw new IllegalArgumentException("wrong length"); + if(rec.size()"+dport+" "+number+"["+data.length+"]"; + return "DATAT "+sport+"->"+dport+" "+number+"["+size+"]"; } @Override @@ -52,8 +55,8 @@ public class DATATPacket extends KLALBPacket { dto.writeInt(sport); dto.writeInt(dport); dto.writeLong(number); - dto.writeChar(data.length); - dto.write(data); + dto.writeChar(size); + dto.write(data,0,size); } @Override @@ -62,7 +65,11 @@ public class DATATPacket extends KLALBPacket { sport=din.readInt(); dport=din.readInt(); number=din.readLong(); - data=new byte[din.readChar()]; - din.readFully(data); + size=din.readChar(); + data=arrayRecycle.create(); + din.readFully(data,0,size); + } + public int getSize() { + return size; } } diff --git a/src/org/kne/cloud/network/klalb/KLALBController.java b/src/org/kne/cloud/network/klalb/KLALBController.java index 46661d3..0b02fce 100644 --- a/src/org/kne/cloud/network/klalb/KLALBController.java +++ b/src/org/kne/cloud/network/klalb/KLALBController.java @@ -264,15 +264,21 @@ public class KLALBController { if (l == null || l.isEmpty()) { throw new NoRouteToHostException("address unreachable: " + addr); } - int count0 = Math.min(count, l.size()); List l2 = (List) ((ArrayList) l).clone(); - LineDecitionComparator ldc = new LineDecitionComparator(l2, priority); + l2.removeAll(packet.getSendRecord()); + if(l2.isEmpty()) { + l2 = (List) ((ArrayList) l).clone(); + } + int count0 = Math.min(count, l2.size()); + LineDecitionComparator ldc = new LineDecitionComparator(l2,packet, priority); Collections.sort(l2,ldc); //System.out.println(l2); for (Iterator iterator = l2.iterator(); iterator.hasNext();) { KLALBRemoteSocket krst = (KLALBRemoteSocket) iterator.next(); krst.sendPacket(packet, priority); + packet.getSendRecord().add(krst); count0--; + Thread.yield(); if (count0 <= 0) break; } diff --git a/src/org/kne/cloud/network/klalb/KLALBPacket.java b/src/org/kne/cloud/network/klalb/KLALBPacket.java index 26d5fd7..8bc5ac0 100644 --- a/src/org/kne/cloud/network/klalb/KLALBPacket.java +++ b/src/org/kne/cloud/network/klalb/KLALBPacket.java @@ -6,6 +6,7 @@ import java.io.Externalizable; import java.io.IOException; import java.io.ObjectInput; import java.io.ObjectOutput; +import java.util.Vector; public abstract class KLALBPacket{ public static final int PING=0; @@ -40,5 +41,11 @@ public abstract class KLALBPacket{ public long getLength() { return 1; } + private Vector sendRecord=new Vector<>(); + + public Vector getSendRecord() { + return sendRecord; + } + } diff --git a/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java b/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java index 1c0b26d..07baf94 100644 --- a/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java +++ b/src/org/kne/cloud/network/klalb/KLALBRemoteSocket.java @@ -35,13 +35,12 @@ public KLALBRemoteSocket(Socket socket) { public KLALBRemoteSocket(Socket socket,Monitor m) { super(socket,m); try { - socket.setSendBufferSize(32768); - socket.setReceiveBufferSize(65535); + //socket.setSendBufferSize(32768); + //socket.setReceiveBufferSize(65535); //System.out.println(socket.getSendBufferSize()+" "+socket.getReceiveBufferSize()); socket.setTcpNoDelay(true); socket.setSoTimeout(10000); } catch (SocketException e1) { - // TODO 自动生成的 catch 块 e1.printStackTrace(); } ThreadTool.makeVDaemonThreadIfSupport("远程接收线程", () -> { @@ -49,20 +48,20 @@ public KLALBRemoteSocket(Socket socket) { while (!isClosed()) { KLALBPacket kpp = getInputStream().readPacket(); - //System.out.println("RX:" + kpp); if(kpp instanceof PINGPacket) { sendPacket(new PONGPacket(((PINGPacket) kpp).getTime()), 65540); }else if(kpp instanceof PONGPacket) { PONGPacket png=(PONGPacket) kpp; getMonitor().updateLatency (System.nanoTime()- png.getTime()); try { - socket.setSoTimeout(1100); + socket.setSoTimeout(600); } catch (SocketException e1) { // TODO 自动生成的 catch 块 e1.printStackTrace(); } }else { + System.out.println("RX:" + kpp); if(kpp instanceof VADDRPacket) { remoteVaddr=((VADDRPacket) kpp).getVaddr(); } @@ -94,15 +93,19 @@ public KLALBRemoteSocket(Socket socket) { PQItem pqi=sendQueue.poll(); if(pqi!=null) { KLALBPacket kpp=pqi.getPacket(); - //if(!(kpp instanceof PONGPacket)) - //System.out.println("TX:" + kpp); + if(!(kpp instanceof PONGPacket)) + System.out.println("TX:" + kpp); if(kpp instanceof PONGPacket) { PONGPacket pp=(PONGPacket) kpp; pp.redeltaTime(System.nanoTime()-pqi.getAddtime()); } getOutputStream().writePacket(kpp); + //System.out.println(sendQueue.size()); }else { - Thread.sleep(1); + synchronized (sendQueue) { + sendQueue.wait(10); + } + //Thread.sleep(1); } } } @@ -117,7 +120,6 @@ public KLALBRemoteSocket(Socket socket) { } } }).start(); - } @@ -178,6 +180,9 @@ public KLALBRemoteSocket(Socket socket) { public void sendPacket(KLALBPacket kp,int priority) { sendQueue.add(new PQItem(kp, priority)); + synchronized (sendQueue) { + sendQueue.notifyAll(); + } } public void setPacketReceiver(Consumer rec) { @@ -194,11 +199,11 @@ public KLALBRemoteSocket(Socket socket) { super.close(); sendQueue.forEach((x)->{ try { + //System.out.println("断线重发"); controller.sendPacketToAddress(remoteVaddr, x.getPacket(), x.getPriority()+1); } catch (NoRouteToHostException e) { } }); - } @Override diff --git a/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java b/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java index 093b1c5..3660b3c 100644 --- a/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java +++ b/src/org/kne/cloud/network/klalb/KLALBVirtualSocketImpl.java @@ -31,12 +31,17 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReferenceArray; import java.util.function.BiConsumer; import java.util.function.Consumer; +import java.util.zip.Deflater; +import java.util.zip.DeflaterOutputStream; +import java.util.zip.Inflater; +import java.util.zip.InflaterInputStream; +import org.kne.cloud.network.mport.VirtualSocketImpl; import org.kne.io.Data; import java.util.*; -public class KLALBVirtualSocketImpl extends SocketImpl { +public class KLALBVirtualSocketImpl extends VirtualSocketImpl { private KLALBController controller; protected int getInputchachesize() { @@ -55,8 +60,19 @@ public class KLALBVirtualSocketImpl extends SocketImpl { this.outputchachesize = outputchachesize; } - private int inputchachesize = 5773 * 100; - private int outputchachesize = 5773 * 100; + private volatile int inputchachesize = LIMIT * 500; + private volatile int outputchachesize = LIMIT * 200; + private volatile boolean nodelay=false; + private volatile long delaytime=1; + + public long getDelaytime() { + return delaytime; + } + + public void setDelaytime(long delaytime) { + this.delaytime = delaytime; + } + private Inet6Address bindaddr;{ try { bindaddr=(Inet6Address) Inet6Address.getByName("::0"); @@ -77,15 +93,20 @@ public class KLALBVirtualSocketImpl extends SocketImpl { try { sendlist.get(i).check(i); } catch (IOException e) { - // TODO 自动生成的 catch 块 e.printStackTrace(); + try { + close(); + } catch (IOException e1) { + e1.printStackTrace(); + } + break; } } } } } - private boolean succeed, refused; + private volatile boolean succeed, refused; protected boolean isListening() { return backlogQueue != null; @@ -93,7 +114,7 @@ public class KLALBVirtualSocketImpl extends SocketImpl { private List inputchache = new ArrayList<>(); private long inputcount = 0; - private boolean avaliable = true; + private volatile boolean avaliable = true; private BiConsumer packReceiver = new BiConsumer() { @Override @@ -104,15 +125,15 @@ public class KLALBVirtualSocketImpl extends SocketImpl { if (backlogQueue.offer(new InetSocketAddress(from, ((SYNTPacket) u).getSport()))) { controller.sendPacketToAddress(from, - new SACKTPacket(localport, ((SYNTPacket) u).getSport()), 65537); + new SACKTPacket(localport, ((SYNTPacket) u).getSport()), 65537,2); } else { controller.sendPacketToAddress(from, new RSTPacket(localport, ((SYNTPacket) u).getSport()), - 65537); + 65537,2); } } else { controller.sendPacketToAddress(from, new RSTPacket(localport, ((SYNTPacket) u).getSport()), - 65537); + 65537,2); } } else if (u instanceof SACKTPacket) { if (connecting) { @@ -148,23 +169,34 @@ public class KLALBVirtualSocketImpl extends SocketImpl { if (kkb == null) break; sendDeque.add(kkb); + synchronized (sendDeque) { + sendDeque.notifyAll(); + + } inputcount++; } } } controller.sendPacketToAddress(from, new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp.getNumber(), - countInputBytes() < inputchachesize), 32768,2); + countInputBytes() < inputchachesize), 32768,1); } else if (u instanceof ACKTPacket) { ACKTPacket ackt = (ACKTPacket) u; avaliable = ackt.isAvaliable(); - AtomicReferencekl=new AtomicReference<>(); + AtomicReferencekl=new AtomicReference<>(); + synchronized (sendlist) { + sendlist.removeIf((tsk)->{ boolean b=tsk.getPacket().getNumber()==ackt.getNumber(); if(b) kl.set(tsk.getPacket()); return b; }); + sendlist.notifyAll(); + } + if(kl.get()!=null) { controller.removeFromSend(from,kl.get()); + DATATPacket.arrayRecycle.recycle(kl.get().getData()); + } } } catch (IOException e) { e.printStackTrace(); @@ -177,14 +209,14 @@ public class KLALBVirtualSocketImpl extends SocketImpl { private int countInputBytes() { AtomicInteger i = new AtomicInteger(0); sendDeque.forEach((c) -> { - i.addAndGet(c.getData().length); + i.addAndGet(c.getSize()); }); return i.get(); } private int countOutputBytes() { AtomicInteger i = new AtomicInteger(0); sendlist.forEach((c) -> { - i.addAndGet(c.getPacket().getData().length); + i.addAndGet(c.getPacket().getSize()); }); return i.get(); } @@ -204,23 +236,34 @@ public class KLALBVirtualSocketImpl extends SocketImpl { @Override public void setOption(int optID, Object value) throws SocketException { - if(optID==SocketOptions.SO_RCVBUF) { + switch(optID) { + case SocketOptions.TCP_NODELAY: + nodelay=(boolean) value; + break; + case SocketOptions.SO_RCVBUF: inputchachesize=(int) value; - }else if(optID==SocketOptions.SO_SNDBUF) { + break; + case SocketOptions.SO_SNDBUF: outputchachesize=(int)value; + break; } } @Override public Object getOption(int optID) throws SocketException { - if(optID==SocketOptions.SO_BINDADDR) { - return bindaddr; - }else if(optID==SocketOptions.SO_RCVBUF) { + switch(optID) { + case SocketOptions.TCP_NODELAY: + return nodelay; + case SocketOptions.SO_RCVBUF: return inputchachesize; - }else if(optID==SocketOptions.SO_SNDBUF) { + case SocketOptions.SO_SNDBUF: return outputchachesize; + case SocketOptions.SO_BINDADDR: + return bindaddr; + default: + return null; } - return null; + } @Override @@ -262,7 +305,6 @@ public class KLALBVirtualSocketImpl extends SocketImpl { } connecting = false; if (succeed) { - } else if (refused) { throw new ConnectException("connect refused"); } else { @@ -329,8 +371,8 @@ public class KLALBVirtualSocketImpl extends SocketImpl { this.acceptedSocketCloseListener = lsr; } - private KVSIInputStream vin; - private KVSIOutputStream vout; + private InputStream vin; + private OutputStream vout; private class KVSIInputStream extends InputStream { private DATATPacket dtp = null; @@ -339,9 +381,7 @@ public class KLALBVirtualSocketImpl extends SocketImpl { private boolean shutdown=false; @Override public int read() throws IOException { - if(shutdown) - return -1; - if (dtp == null || (count >= dtp.getData().length && dtp.getData().length != 0)) { + if (dtp == null ) { count = 0; while (true) { if (isClosed()) @@ -356,16 +396,24 @@ public class KLALBVirtualSocketImpl extends SocketImpl { break; } try { - Thread.sleep(1); + synchronized (sendDeque) { + sendDeque.wait(100); + } } catch (InterruptedException e) { e.printStackTrace(); } } } - if (dtp.getData().length == 0) { + if (dtp.getSize() == 0) { return -1; - } else - return dtp.getData()[count++] & 0xff; + } else { + int ret= dtp.getData()[count++] & 0xff; + if(count==dtp.getSize()) { + DATATPacket.arrayRecycle.recycle(dtp.getData()); + dtp=null; + } + return ret; + } } @Override @@ -379,21 +427,78 @@ public class KLALBVirtualSocketImpl extends SocketImpl { } len = Math.min(len, available()); - - int c = read(); - if (c == -1) { - return -1; + + + if (dtp == null ) { + count = 0; + while (true) { + if (isClosed()) + throw new SocketException("Socket is closed"); + DATATPacket dtp2 = sendDeque.poll(); + if (dtp2 != null) { + dtp = dtp2; + if(countInputBytes() < inputchachesize) { + controller.sendPacketToAddress((Inet6Address) address, new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp2.getNumber(), + true), 32768); + } + break; + } + try { + synchronized (sendDeque) { + sendDeque.wait(100); + + } + } catch (InterruptedException e) { + e.printStackTrace(); + } + } + } + if (dtp.getSize() == 0) { + return -1; + } else { + b[off]= dtp.getData()[count++] ; + if(count==dtp.getSize()) { + DATATPacket.arrayRecycle.recycle(dtp.getData()); + dtp=null; + } } - b[off] = (byte) c; - int i = 1; try { for (; i < len; i++) { - c = read(); - if (c == -1) { - break; + + if (dtp == null ) { + count = 0; + while (true) { + if (isClosed()) + throw new SocketException("Socket is closed"); + DATATPacket dtp2 = sendDeque.poll(); + if (dtp2 != null) { + dtp = dtp2; + if(countInputBytes() < inputchachesize) { + controller.sendPacketToAddress((Inet6Address) address, new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp2.getNumber(), + true), 32768); + } + break; + } + try { + synchronized (sendDeque) { + sendDeque.wait(100); + + } + } catch (InterruptedException e) { + e.printStackTrace(); + } + } + } + if (dtp.getSize() == 0) { + break; + } else { + b[off + i]= dtp.getData()[count++] ; + if(count==dtp.getSize()) { + DATATPacket.arrayRecycle.recycle(dtp.getData()); + dtp=null; + } } - b[off + i] = (byte) c; } } catch (IOException ee) { } @@ -402,40 +507,92 @@ public class KLALBVirtualSocketImpl extends SocketImpl { @Override public void close() throws IOException { - shutdown=true; + } @Override public int available() throws IOException { AtomicInteger i = new AtomicInteger(0); sendDeque.forEach((V) -> { - i.addAndGet(V.getData().length); + i.addAndGet(V.getSize()); }); if (dtp != null) - i.addAndGet(dtp.getData().length - count); + i.addAndGet(dtp.getSize() - count); return i.get(); } } - private long outputcount = 0; + private volatile long outputcount = 0; + private static final int LIMIT=65535; private class KVSIOutputStream extends OutputStream { - - private byte[] cache=new byte[65535]; + private byte[] cache=DATATPacket.arrayRecycle.create(); private int count=0; + private Object lock=new Object(); @Override public void write(int b) throws IOException { if (isClosed()) throw new SocketException("Socket is closed"); + synchronized (lock) { + cache[count++]=(byte) b; - if (count >=5773 ) {//1429? + if (count >= LIMIT) {//1429?5773?8669 + flush0(); + }else { flush(); + } } } + + @Override + public void write(byte[] b, int off, int len) throws IOException { + if (isClosed()) + throw new SocketException("Socket is closed"); + int ol=off+len; + synchronized (lock) { + for (int i = off; i < ol; i++) { + cache[count++]=b[i]; + if (count >= LIMIT) {//1429?5773?8669 + flush0(); + } + } + flush(); + } + } + + private volatile TimerTask tt; @Override public void flush() throws IOException { + if(nodelay) { + synchronized (lock) { + flush0(); + } + }else { + if(tt==null) { + tt=new TimerTask() { + + @Override + public void run() { + if(isClosed()) + cancel(); + try { + synchronized (lock) { + flush0(); + } + } catch (IOException e) { + e.printStackTrace(); + } + } + }; + new Timer("粘包计时线程").scheduleAtFixedRate(tt, 0, delaytime); + } + } + } + + + private void flush0() throws IOException { if (count > 0) { try { while (!avaliable) { @@ -446,20 +603,18 @@ public class KLALBVirtualSocketImpl extends SocketImpl { } try { while(countOutputBytes()>outputchachesize) { - Thread.sleep(1); + synchronized (sendlist) { + sendlist.wait(10); + } } } catch (InterruptedException e) { // TODO 自动生成的 catch 块 e.printStackTrace(); } - byte[]ba; - if(count==cache.length) { - ba=cache; - cache=new byte[cache.length]; - }else { - ba=Arrays.copyOf(cache, count); - } - SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, ba), 5,10); + byte[] ba=cache; + cache=DATATPacket.arrayRecycle.create(); + + SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, ba,count), 5,20); x.run(); sendlist.add(x); /*controller.sendPacketToAddress((Inet6Address) address, @@ -470,28 +625,36 @@ public class KLALBVirtualSocketImpl extends SocketImpl { @Override public void close() throws IOException { - flush(); - SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, new byte[0]), 5,10); + flush0(); + SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, DATATPacket.arrayRecycle.create(),0), 5,20); x.run(); sendlist.add(x); /*controller.sendPacketToAddress((Inet6Address) address, new DATATPacket(localport, port, outputcount++, new byte[0]), 5);*/ + tt.cancel(); } } @Override - protected InputStream getInputStream() throws IOException { + public InputStream getInputStream() throws IOException { if (vin == null) { vin = new KVSIInputStream(); + int val= vin.read(); + if(val==1) { + vin=new InflaterInputStream(vin,new Inflater(true),LIMIT); + } } return vin; } @Override - protected OutputStream getOutputStream() throws IOException { + public OutputStream getOutputStream() throws IOException { if (vout == null) { vout = new KVSIOutputStream(); + vout.write(1); + vout.flush(); + vout=new DeflaterOutputStream(vout, new Deflater(Deflater.BEST_COMPRESSION, true), LIMIT, true); } return vout; } @@ -521,17 +684,23 @@ public class KLALBVirtualSocketImpl extends SocketImpl { protected void close() throws IOException { if (!isListening() && !isClosed()) { closed = true; + try { + controller.sendPacketToAddress((Inet6Address) super.address, new RSTPacket(super.localport, super.port), - 65537); + 65537,2); + } catch (NoRouteToHostException e) { + } + if (acceptedSocketCloseListener != null) + acceptedSocketCloseListener.accept(this); + else + controller.unbind(this); } - if (acceptedSocketCloseListener != null) - acceptedSocketCloseListener.accept(this); - else - controller.unbind(this); sendCheckTask.cancel(); + +//new Exception().printStackTrace(); } - private boolean closed; + private volatile boolean closed; private boolean isClosed() { return closed; diff --git a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java index 263a909..2999198 100644 --- a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java +++ b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java @@ -8,38 +8,42 @@ import java.util.Map; import java.util.concurrent.atomic.AtomicLong; public class LineDecitionComparator implements Comparator { - private Map predictedlatencys=new HashMap<>(); - public LineDecitionComparator (List krs,int priority) { + private Map predictedlatencys=new HashMap<>(); + public LineDecitionComparator (List krs,KLALBPacket curr,int priority) { for (Iterator iterator = krs.iterator(); iterator.hasNext();) { KLALBRemoteSocket klalbRemoteSocket = (KLALBRemoteSocket) iterator.next(); - long x=klalbRemoteSocket.getMonitor().getLatency()>>1; - x+=klalbRemoteSocket.getMonitor().getJitter()>>1; + double x=klalbRemoteSocket.getMonitor().getLatency()/2; + x+=klalbRemoteSocket.getMonitor().getJitter()/2; AtomicLong al=new AtomicLong(0); klalbRemoteSocket.getSendQueue().forEach((v)->{ if(v.getPriority()>=priority) { al.addAndGet(v.getPacket().getLength()); } }); - long speed=klalbRemoteSocket.getMonitor(). getOutSpeedMax(); + al.addAndGet(curr.getLength()); + long speed=klalbRemoteSocket.getMonitor(). getOutSpeed(); if(speed==0) { if(al.get()>0) { x=Long.MAX_VALUE; } }else { - x+=al.get()*1000000000/speed; + x+=al.get()*1000000000.0/speed; } - //System.out.println(x); + + /*System.out.println(klalbRemoteSocket); + System.out.println(klalbRemoteSocket.getSendQueue().size()); + System.out.println(x);*/ predictedlatencys.put(klalbRemoteSocket, x); } } - public Map getPredictedlatencys() { + public Map getPredictedlatencys() { return predictedlatencys; } @Override public int compare(KLALBRemoteSocket o1, KLALBRemoteSocket o2) { - long t1=predictedlatencys.get(o1); - long t2=predictedlatencys.get(o2); + double t1=predictedlatencys.get(o1); + double t2=predictedlatencys.get(o2); if(t1>t2) { return 1; }else if(t1{ } } public synchronized void addHostPort(MultipurposeSocketAddress v) { - LineEntry le=new LineEntry(v); if(putIfAbsent(v,le)==null) le.open(); @@ -74,10 +73,11 @@ public class LineManager extends Hashtable{ private MultipurposeSocketAddress hp; private KLALBRemoteSocket krs; private volatile boolean flag=true; - private volatile boolean online=false; private Monitor monitor=new Monitor(); - public boolean isOnline() { - return online; + + + public Monitor getMonitor() { + return monitor; } public MultipurposeSocketAddress getHp() { @@ -90,7 +90,6 @@ public class LineManager extends Hashtable{ } public String toString() { StringBuilder sb=new StringBuilder(); - sb.append(online?"●在线 ":"○离线 "); sb.append(hp.toString()); sb.append('\t'); sb.append(monitor.toString()); @@ -102,7 +101,6 @@ public class LineManager extends Hashtable{ sb.append(hp.toString()); sb.append('\n'); - sb.append(online?"●在线\t":"○离线\t"); sb.append(monitor.toString2()); return sb.toString(); } @@ -111,24 +109,34 @@ public class LineManager extends Hashtable{ public void run() { while(true){ try { - krs=new KLALBRemoteSocket(hp.connectSocket(),monitor); + monitor.setState(Monitor.CONNECTING); + krs=new KLALBRemoteSocket(hp.connectSocket(3000),monitor); //System.out.println("open"); controller.addRemoteSocket(krs); - online=true; + monitor.resetCoolingTime(); + monitor.setState(Monitor.ONLINE); krs.waitforlose(); //System.out.println("close"); - online=false; } catch (UnknownHostException e) { } catch (IOException e) { //e.printStackTrace(); + }finally { + monitor.setState(Monitor.OFFLINE); + try { + if(krs!=null) + krs.close(); + } catch (IOException e) { + e.printStackTrace(); + } } if(!flag) { break; } try { synchronized (lock) { - lock.wait(10000); + lock.wait(monitor.getCoolingTime()); } + monitor.incCoolingTime(); } catch (InterruptedException e) { e.printStackTrace(); } diff --git a/src/org/kne/cloud/network/klalb/Monitor.java b/src/org/kne/cloud/network/klalb/Monitor.java index cf2eda3..7a027c5 100644 --- a/src/org/kne/cloud/network/klalb/Monitor.java +++ b/src/org/kne/cloud/network/klalb/Monitor.java @@ -1,9 +1,35 @@ package org.kne.cloud.network.klalb; +import java.util.Timer; +import java.util.TimerTask; import java.util.concurrent.atomic.AtomicLong; +import java.util.function.Consumer; + import static org.kne.cloud.network.klalb.KLALBUtils.*; public class Monitor { + private static Timer t=new Timer("带宽测量线程",true); + + + private TimerTask ptt=new TimerTask() { + + @Override + public void run() { + runMonitor(); + } + }; + protected void runMonitor() { + updateSpeed(); + + } + @Override + protected void finalize() throws Throwable { + ptt.cancel(); + } + public Monitor() { + + t.scheduleAtFixedRate(ptt, 1000, 1000); + } private AtomicLong inTraffic=new AtomicLong(); private AtomicLong outTraffic=new AtomicLong(); @@ -19,21 +45,34 @@ public class Monitor { private volatile long latency ; private volatile long jitter ; + + private volatile int state=0; + public static final int OFFLINE=0; + public static final int CONNECTING=1; + public static final int ONLINE=2; + private volatile long coolingTime; + { + resetCoolingTime(); + } + + public int getState() { + return state; + } + public void setState(int state) { + this.state = state; + } public long getInSpeedMax() { return inSpeedMax; } public void updateLatency(long newlatency) { - long vj=Math.abs( latency-newlatency);; - if(jitter>1; + //} } public long getOutSpeedMax() { return outSpeedMax; @@ -67,40 +106,98 @@ public class Monitor { return outSpeed; } + private long updatetime=System.nanoTime(); public void updateSpeed() { + long d=System.nanoTime()-updatetime; long i=inTraffic.get()-inTrafficOld.get(); long o=outTraffic.get()-outTrafficOld.get(); inTrafficOld.set(inTraffic.get()); outTrafficOld.set(outTraffic.get()); - - inSpeed= i; - outSpeed= o; + + inSpeed= i*1000000000/d; + outSpeed= o*1000000000/d; + updatetime=System.nanoTime(); + //System.out.println(d+" "+o+" "+i); if(inSpeed>inSpeedMax) { inSpeedMax=inSpeed; }else { - inSpeedMax=(inSpeedMax*999+inSpeed)/1000; + inSpeedMax=(inSpeedMax*9999+inSpeed)/10000; } if(outSpeed>outSpeedMax) { outSpeedMax=outSpeed; }else { - outSpeedMax=(outSpeedMax*999+outSpeed)/1000; + outSpeedMax=(outSpeedMax*9999+outSpeed)/10000; + } + + if(changeListener!=null) { + changeListener.accept(this); } } + private ConsumerchangeListener; + public Consumer getChangeListener() { + return changeListener; + } + public void setChangeListener(Consumer changeListener) { + this.changeListener = changeListener; + } public long getLatency() { return latency; } public String toString() { StringBuilder sb=new StringBuilder(); + sb.append(getDsc()); + sb.append('\t'); sb.append(bytesUnit(outTraffic.get())).append("\u2191\t").append(bytesUnit(inTraffic.get())).append("\u2193\t").append(bytesUnit(outSpeed)).append("/s\u2191\t").append(bytesUnit(inSpeed)).append("/s\u2193\t").append(latency/1000000).append("ms"); return sb.toString(); } public String toString2() { StringBuilder sb=new StringBuilder(); - sb.append(bytesUnit(outTraffic.get())).append("\t").append(bytesUnit(inTraffic.get())).append("\t").append(bytesUnit(outSpeed)).append("/s\t").append(bytesUnit(inSpeed)).append("/s\t").append(latency/1000000).append("ms").append("\t").append(jitter/1000000).append("ms"); - return sb.toString(); + sb.append(getDsc()); + sb.append('\t'); + sb.append(bytesUnit(outTraffic.get())).append("\t").append(bytesUnit(inTraffic.get())).append("\t").append(bytesUnit(outSpeed)).append("/s\t").append(bytesUnit(inSpeed)).append("/s\t").append(latency/1000000).append("ms").append("\t").append(jitter/1000000).append("ms\t"); + + if(state==OFFLINE) { + sb.append((coolingTime-(System.currentTimeMillis()-mls))/1000L); + sb.append('s'); + }else { + sb.append('/'); + } + + sb.append('\n'); + sb.append(bytesUnit(outSpeedMax)); + + sb.append('\n'); + sb.append(bytesUnit(inSpeedMax)); + + return sb.toString(); + } + private String getDsc() { + switch(state) { + case OFFLINE: + return "○离线"; + case CONNECTING: + return "◐连接中"; + case ONLINE: + return "●在线"; + } + return null; + } + private volatile long mls; + public long getCoolingTime() { + mls=System.currentTimeMillis(); + return coolingTime; + } + public void incCoolingTime() { + coolingTime<<=1; + if(coolingTime>300000) { + coolingTime=300000; + } + } + public void resetCoolingTime() { + coolingTime=10000; } } diff --git a/src/org/kne/cloud/network/klalb/MonitoredSocket.java b/src/org/kne/cloud/network/klalb/MonitoredSocket.java index a19e8d8..8cf78f4 100644 --- a/src/org/kne/cloud/network/klalb/MonitoredSocket.java +++ b/src/org/kne/cloud/network/klalb/MonitoredSocket.java @@ -17,27 +17,20 @@ public class MonitoredSocket extends FilterSocket { - private static Timer t=new Timer("带宽测量线程",true); private Monitor monitor; - private TimerTask ptt=new TimerTask() { - - @Override - public void run() { - runMonitor(); - } - }; + + @Override + public synchronized void close() throws IOException { + super.close(); + } public MonitoredSocket(Socket socket) { this(socket,new Monitor()); } public MonitoredSocket(Socket socket,Monitor monitor) { super(socket); - t.scheduleAtFixedRate(ptt, 1000, 1000); this.monitor=monitor; } - protected void runMonitor() { - monitor.updateSpeed(); - - } + diff --git a/src/org/kne/cloud/network/klalb/SendTask.java b/src/org/kne/cloud/network/klalb/SendTask.java index 5d078e6..e738a96 100644 --- a/src/org/kne/cloud/network/klalb/SendTask.java +++ b/src/org/kne/cloud/network/klalb/SendTask.java @@ -56,8 +56,10 @@ public class SendTask { + public void check(int number) throws IOException { - long limit= (1<{ KLALBVirtualSocket kvs=null; try { + s.setTcpNoDelay(true); kvs=new KLALBVirtualSocket(kc, kr.getRemoteVaddr(), 23333); - new SocketBridge(s, kvs).run(); + kvs.setTcpNoDelay(true); + SocketBridge dsb=new SocketBridge(s, kvs); + dsb.getSab().setDelay(1); + dsb.run(); + /*CompressedSocketBridge csb=new CompressedSocketBridge(s, kvs); + csb.getSab().setDelay(1); + csb.run();*/ } catch (IOException e) { e.printStackTrace(); }finally { @@ -56,7 +65,7 @@ public static void main(String[] args) throws UnknownHostException, IOException String s=scn.nextLine(); switch(s) { case "state": - System.out.println("状态\t上传流量\t下载流量\t上传速度\t下载速度\t延迟\t抖动"); + System.out.println("状态\t上传流量\t下载流量\t上传速度\t下载速度\t延迟\t抖动\t下一次重试"); LineManager le=kc.getLineManager(); for (Iterator> iterator = le.entrySet().iterator(); iterator.hasNext();) { java.util.Map.Entry hostPort = iterator.next(); diff --git a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java index 8a2544b..8027082 100644 --- a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java +++ b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java @@ -119,7 +119,14 @@ public class SimpleKLALBServer { Socket s=null; try { s=mpsa.connectSocket(); - new SocketBridge(soc, s).run(); + s.setTcpNoDelay(true); + soc.setTcpNoDelay(true); + SocketBridge dsb=new SocketBridge(s, soc); + dsb.getSab().setDelay(1); + dsb.run(); + /*CompressedSocketBridge csb=new CompressedSocketBridge(s, soc); + csb.getSab().setDelay(1); + csb.run();*/ } catch (UnknownHostException e) { // TODO 自动生成的 catch 块 e.printStackTrace(); diff --git a/src/org/kne/cloud/network/mport/SocketBridge.java b/src/org/kne/cloud/network/mport/SocketBridge.java index 3aac927..e4cedc2 100644 --- a/src/org/kne/cloud/network/mport/SocketBridge.java +++ b/src/org/kne/cloud/network/mport/SocketBridge.java @@ -1,25 +1,35 @@ package org.kne.cloud.network.mport; import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; import java.net.*; import org.kne.io.Task; public class SocketBridge extends Task{ private Socket a; private Socket b; - public SocketBridge(Socket a, Socket b) { + public StreamBridge getSab() { + return sab; + } + public StreamBridge getSba() { + return sba; + } + private StreamBridge sab; + private StreamBridge sba; + public SocketBridge(Socket a, Socket b) throws IOException { super(); this.a = a; this.b = b; + sab = new StreamBridge(getAIN(), getBOUT()); + sba = new StreamBridge(getBIN(), getAOUT()); } @Override protected void runTask() { try { - StreamBridge sba=new StreamBridge(a.getInputStream(), b.getOutputStream()); - StreamBridge sbb=new StreamBridge(b.getInputStream(), a.getOutputStream()); - sba.runAtNewThread("SocketBridge A->B thread"); - sbb.runAtNewThread("SocketBridge B->A thread"); + sab.runAtNewThread("SocketBridge A->B thread"); + sba.runAtNewThread("SocketBridge B->A thread"); + sab.waitfortask(); sba.waitfortask(); - sbb.waitfortask(); } catch (Exception e) { e.printStackTrace(); }finally { @@ -37,5 +47,23 @@ public class SocketBridge extends Task{ } } } + public Socket getA() { + return a; + } + public Socket getB() { + return b; + } + protected OutputStream getBOUT() throws IOException { + return b.getOutputStream(); + } + protected InputStream getBIN() throws IOException { + return b.getInputStream(); + } + protected OutputStream getAOUT() throws IOException { + return a.getOutputStream(); + } + protected InputStream getAIN() throws IOException { + return a.getInputStream(); + } } diff --git a/src/org/kne/cloud/network/mport/StreamBridge.java b/src/org/kne/cloud/network/mport/StreamBridge.java index c049ef0..73f1c7c 100644 --- a/src/org/kne/cloud/network/mport/StreamBridge.java +++ b/src/org/kne/cloud/network/mport/StreamBridge.java @@ -9,6 +9,7 @@ public class StreamBridge extends Task{ private InputStream in; private OutputStream out; private int blocksize=65535; + private long delay=0; public StreamBridge(InputStream in, OutputStream out) { super(); Objects.requireNonNull(in); @@ -25,6 +26,11 @@ public class StreamBridge extends Task{ while ((v=in.read(b))!=-1) { out.write(b,0,v); out.flush(); + try { + Thread.sleep(delay); + } catch (InterruptedException e) { + e.printStackTrace(); + } } } catch (IOException e) { // TODO 自动生成的 catch 块 @@ -51,6 +57,14 @@ public class StreamBridge extends Task{ public void setBlocksize(int blocksize) { this.blocksize = blocksize; } + public long getDelay() { + return delay; + } + + public void setDelay(long delay) { + this.delay = delay; + } + public void runAtNewThread() { runAtNewThread("StreamBridge thread"); } diff --git a/src/org/kne/cloud/network/mport/ThreadTool.java b/src/org/kne/cloud/network/mport/ThreadTool.java index a02920d..2982542 100644 --- a/src/org/kne/cloud/network/mport/ThreadTool.java +++ b/src/org/kne/cloud/network/mport/ThreadTool.java @@ -21,7 +21,7 @@ public class ThreadTool { }catch(Throwable e) { //e.printStackTrace(); if(first) { - System.out.println("请使用java19以上版本以提高性能!"); + System.out.println("请使用java19以上版本并开启--enable-preview选项以提高性能!"); first=false; } return new Thread(r, name); diff --git a/src/org/kne/cloud/network/mport/VirtualServerSocket.java b/src/org/kne/cloud/network/mport/VirtualServerSocket.java index 0d2bf88..0fab1ce 100644 --- a/src/org/kne/cloud/network/mport/VirtualServerSocket.java +++ b/src/org/kne/cloud/network/mport/VirtualServerSocket.java @@ -14,8 +14,8 @@ import javax.sql.rowset.RowSetMetaDataImpl; public abstract class VirtualServerSocket extends ServerSocket { public VirtualServerSocket(SocketImpl si) throws IOException { - super(); - try { + super(si); + /*try { Class c=ServerSocket.class; Field con =c.getDeclaredField("impl"); con.setAccessible(true); @@ -42,7 +42,7 @@ public abstract class VirtualServerSocket extends ServerSocket { } catch (InvocationTargetException e) { // TODO 自动生成的 catch 块 e.printStackTrace(); - } + }*/ } @Override diff --git a/src/org/kne/cloud/network/mport/VirtualSocket.java b/src/org/kne/cloud/network/mport/VirtualSocket.java index ba5c088..25ef25c 100644 --- a/src/org/kne/cloud/network/mport/VirtualSocket.java +++ b/src/org/kne/cloud/network/mport/VirtualSocket.java @@ -13,8 +13,21 @@ import java.nio.channels.SocketChannel; public abstract class VirtualSocket extends Socket { - public VirtualSocket(SocketImpl si) throws SocketException { + private VirtualSocketImpl si; + + public VirtualSocket(VirtualSocketImpl si) throws SocketException { super(si); + this.si=si; + } + + @Override + public InputStream getInputStream() throws IOException { + return si.getInputStream(); + } + + @Override + public OutputStream getOutputStream() throws IOException { + return si.getOutputStream(); } diff --git a/src/org/kne/cloud/network/mport/VirtualSocketImpl.java b/src/org/kne/cloud/network/mport/VirtualSocketImpl.java new file mode 100644 index 0000000..44a0097 --- /dev/null +++ b/src/org/kne/cloud/network/mport/VirtualSocketImpl.java @@ -0,0 +1,19 @@ +package org.kne.cloud.network.mport; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.net.InetAddress; +import java.net.SocketAddress; +import java.net.SocketException; +import java.net.SocketImpl; + +public abstract class VirtualSocketImpl extends SocketImpl { + + @Override + public abstract InputStream getInputStream() throws IOException ; + + @Override + public abstract OutputStream getOutputStream() throws IOException; + +} diff --git a/src/org/kne/cloud/network/nathole/NatholeTestC.java b/src/org/kne/cloud/network/nathole/NatholeTestC.java new file mode 100644 index 0000000..c0b2801 --- /dev/null +++ b/src/org/kne/cloud/network/nathole/NatholeTestC.java @@ -0,0 +1,32 @@ +package org.kne.cloud.network.nathole; + +import java.io.IOException; +import java.net.BindException; +import java.net.InetSocketAddress; +import java.net.Socket; +import java.net.SocketAddress; +import java.net.SocketOption; +import java.net.SocketOptions; +import java.nio.channels.SocketChannel; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +public class NatholeTestC { + public static ExecutorService exc=Executors.newCachedThreadPool(); + public static void main(String[] args) throws IOException { + for (int i1 = 1024; i1 < 65536; i1++) { + int i2=i1; + exc.execute(()->{ + try { + Socket s=new Socket(); + System.out.println(i2); + s.bind(new InetSocketAddress("0.0.0.0", i2)); + s.connect(new InetSocketAddress("10.235.16.1", 10400), 100); + System.out.println("成功!"); + System.exit(0); + } catch (IOException e) { + } + }); + } + } +} diff --git a/src/org/kne/cloud/network/nathole/NatholeTestS.java b/src/org/kne/cloud/network/nathole/NatholeTestS.java new file mode 100644 index 0000000..e08c3bf --- /dev/null +++ b/src/org/kne/cloud/network/nathole/NatholeTestS.java @@ -0,0 +1,36 @@ +package org.kne.cloud.network.nathole; + +import java.io.IOException; +import java.net.BindException; +import java.net.InetSocketAddress; +import java.net.ServerSocket; +import java.net.Socket; +import java.nio.channels.SocketChannel; + +public class NatholeTestS { + public static void main(String[] args) throws IOException { + for (int i1 = 1024; i1 < 65536; i1++) { + try { + SocketChannel sc= SocketChannel.open(); + sc.configureBlocking(false) ; + sc.bind(new InetSocketAddress("0.0.0.0", i1)); + sc.connect(new InetSocketAddress("10.235.16.1",10300)); + sc.close(); + int i2=i1; + new Thread(()->{ + try { + ServerSocket ssk=new ServerSocket(i2); + while(true) { + Socket s=ssk.accept(); + System.out.println(s); + } + } catch (IOException e) { + e.printStackTrace(); + } + } ).start(); + System.out.println(i1); + }catch(BindException e) { + } + } + } +}