From b4fccf3a03f171a0182c68045444942685664400 Mon Sep 17 00:00:00 2001 From: Administrator Date: Tue, 30 Jan 2024 16:54:09 +0800 Subject: [PATCH] KLALB UDP TEST --- .classpath | 7 +- KLALB协议规范V2.0.docx | Bin 19427 -> 19732 bytes kserver.ini | 4 +- linetable - 副本.txt | 3 +- linetable.txt | 22 +- .../DefaultServerSocketFactory.java | 8 +- .../{mport => }/DefaultSocketFactory.java | 6 +- .../{mport => }/FilterServerSocket.java | 2 +- .../network/{mport => }/FilterSocket.java | 2 +- .../network/{mport => }/HTTPDetectorItem.java | 2 +- .../{mport => }/HTTPSDetectorItem.java | 2 +- .../network/{mport => }/HostPortMap.java | 2 +- .../{mport => }/KLALBDetectorItem.java | 2 +- .../MultipurposeSocketAddress.java | 103 +- .../network/{mport => }/NetworkService.java | 2 +- .../network/{mport => }/PortMultiUse.java | 2 +- .../cloud/network/{mport => }/PortRelay.java | 9 +- .../cloud/network/{mport => }/Protocol.java | 2 +- .../network/{mport => }/ProtocolDetector.java | 2 +- .../{mport => }/ProtocolDetectorItem.java | 2 +- .../ProtocolDetectorServerSocket.java | 2 +- .../ProtocolDetectorServerSocketFactory.java | 10 +- .../{mport => }/ProtocolDetectorSocket.java | 2 +- .../network/{mport => }/ProtocolStack.java | 2 +- .../ProxyProfileEntry.java => Proxy.java} | 10 +- .../{mport => }/ProxyProfileAnalyser.java | 0 .../{mport => }/ProxyProfileExecutor.java | 0 .../network/{mport => }/RDPDetectorItem.java | 2 +- .../network/{mport => }/SSHDetectorItem.java | 2 +- .../network/{mport => }/ServiceElement.java | 2 +- ...leEntry.java => ServiceToSocketProxy.java} | 12 +- .../network/{mport => }/SocketBridge.java | 38 +- .../TCPListener.java => SocketListener.java} | 44 +- ...leEntry.java => SocketToServiceProxy.java} | 12 +- .../cloud/network/SocketToSocketProxy.java | 130 ++ .../network/{mport => }/StreamBridge.java | 23 +- .../cloud/network/{mport => }/ThreadTool.java | 18 +- .../{mport => }/VirtualServerSocket.java | 11 +- .../network/{mport => }/VirtualSocket.java | 14 +- .../{mport => }/VirtualSocketImpl.java | 2 +- .../kne/cloud/network/klalb/ACKTPacket.java | 70 +- .../{LINESPacket.java => ADDLINESPacket.java} | 12 +- .../cloud/network/klalb/ByteArrayRecycle.java | 4 +- .../network/klalb/CompressedSocketBridge.java | 14 +- .../kne/cloud/network/klalb/DATATPacket.java | 109 +- .../cloud/network/klalb/KLALBController.java | 483 +++---- .../cloud/network/klalb/KLALBInputStream.java | 49 +- .../kne/cloud/network/klalb/KLALBMain.java | 76 +- .../network/klalb/KLALBOutputStream.java | 12 +- .../kne/cloud/network/klalb/KLALBPacket.java | 30 +- .../kne/cloud/network/klalb/KLALBUtils.java | 21 + .../klalb/KLALBVirtualServerSocket.java | 4 +- .../KLALBVirtualServerSocketFactory.java | 5 + .../network/klalb/KLALBVirtualSocket.java | 16 +- .../klalb/KLALBVirtualSocketFactory.java | 4 + .../network/klalb/KLALBVirtualSocketImpl.java | 1105 ++++++++++++----- .../network/klalb/LineDecitionComparator.java | 38 +- src/org/kne/cloud/network/klalb/Monitor.java | 104 +- .../cloud/network/klalb/MonitoredSocket.java | 2 +- .../kne/cloud/network/klalb/PONGPacket.java | 36 +- .../kne/cloud/network/klalb/RSTPacket.java | 2 +- src/org/kne/cloud/network/klalb/SendTask.java | 7 +- .../network/klalb/SimpleKLALBClient.java | 58 +- .../network/klalb/SimpleKLALBServer.java | 98 +- .../SocketToSocketProxyProfileEntry.java | 16 - .../cloud/network/nathole/NatholeTestS.java | 4 +- src/org/kne/debug/TimeDebugger.java | 7 +- 67 files changed, 1977 insertions(+), 929 deletions(-) rename src/org/kne/cloud/network/{mport => }/DefaultServerSocketFactory.java (78%) rename src/org/kne/cloud/network/{mport => }/DefaultSocketFactory.java (88%) rename src/org/kne/cloud/network/{mport => }/FilterServerSocket.java (95%) rename src/org/kne/cloud/network/{mport => }/FilterSocket.java (95%) rename src/org/kne/cloud/network/{mport => }/HTTPDetectorItem.java (91%) rename src/org/kne/cloud/network/{mport => }/HTTPSDetectorItem.java (86%) rename src/org/kne/cloud/network/{mport => }/HostPortMap.java (87%) rename src/org/kne/cloud/network/{mport => }/KLALBDetectorItem.java (87%) rename src/org/kne/cloud/network/{mport => }/MultipurposeSocketAddress.java (50%) rename src/org/kne/cloud/network/{mport => }/NetworkService.java (84%) rename src/org/kne/cloud/network/{mport => }/PortMultiUse.java (94%) rename src/org/kne/cloud/network/{mport => }/PortRelay.java (78%) rename src/org/kne/cloud/network/{mport => }/Protocol.java (87%) rename src/org/kne/cloud/network/{mport => }/ProtocolDetector.java (94%) rename src/org/kne/cloud/network/{mport => }/ProtocolDetectorItem.java (76%) rename src/org/kne/cloud/network/{mport => }/ProtocolDetectorServerSocket.java (94%) rename src/org/kne/cloud/network/{mport => }/ProtocolDetectorServerSocketFactory.java (82%) rename src/org/kne/cloud/network/{mport => }/ProtocolDetectorSocket.java (90%) rename src/org/kne/cloud/network/{mport => }/ProtocolStack.java (65%) rename src/org/kne/cloud/network/{mport/ProxyProfileEntry.java => Proxy.java} (65%) rename src/org/kne/cloud/network/{mport => }/ProxyProfileAnalyser.java (100%) rename src/org/kne/cloud/network/{mport => }/ProxyProfileExecutor.java (100%) rename src/org/kne/cloud/network/{mport => }/RDPDetectorItem.java (86%) rename src/org/kne/cloud/network/{mport => }/SSHDetectorItem.java (87%) rename src/org/kne/cloud/network/{mport => }/ServiceElement.java (91%) rename src/org/kne/cloud/network/{mport/ServiceToSocketProxyProfileEntry.java => ServiceToSocketProxy.java} (53%) rename src/org/kne/cloud/network/{mport => }/SocketBridge.java (55%) rename src/org/kne/cloud/network/{mport/TCPListener.java => SocketListener.java} (52%) rename src/org/kne/cloud/network/{mport/SocketToServiceProxyProfileEntry.java => SocketToServiceProxy.java} (55%) create mode 100644 src/org/kne/cloud/network/SocketToSocketProxy.java rename src/org/kne/cloud/network/{mport => }/StreamBridge.java (71%) rename src/org/kne/cloud/network/{mport => }/ThreadTool.java (72%) rename src/org/kne/cloud/network/{mport => }/VirtualServerSocket.java (82%) rename src/org/kne/cloud/network/{mport => }/VirtualSocket.java (68%) rename src/org/kne/cloud/network/{mport => }/VirtualSocketImpl.java (88%) rename src/org/kne/cloud/network/klalb/{LINESPacket.java => ADDLINESPacket.java} (81%) delete mode 100644 src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java diff --git a/.classpath b/.classpath index 0039a3b..407aa40 100644 --- a/.classpath +++ b/.classpath @@ -2,8 +2,13 @@ + + + + + - + diff --git a/KLALB协议规范V2.0.docx b/KLALB协议规范V2.0.docx index 0886dc081bf071ea9b7964be24c03d05f788f572..e36c7b1c4bdf84f9c8717a75349b83b06af4540f 100644 GIT binary patch delta 6832 zcmZ9RRa6uJw}ppBy1N@hx`rH@p+h>PLAr*9p<8N@2I-Ve32Bg$E{E<8L6E-w_dfh< z-Iw!l&N}<7^SJl7lYx+yjgTyeqT{5dXObdH4b^fD;~|RN$T*e4TQ5Wz0_yPK9LzOl z#azz_>>*ATIRA3E*{Dm9UpZMn2ewJv|B_A5sQvS^06T->Pchyj4LXYo*YW%YX^#1M zz_7N~tR9Ba^II}#mV{(-O*%&Np1C+Fz3Q>>Vik%^l;E@$VMmKwj`DG$5i8`AwXQp= zBXm@q8a>`J!492^MDC4EflKI1a=QrpH7Y&zPvCx2I|o{B3Zm$OHs(q|(p*$&GqtpE zhSNuEmUwguRFC8i0df|5GNpA*wJgZ%d=8`hdb^u&9uSKQJZGe{@RZ3YkMZ%N@g>f~ zoR743LTj?v>er$ggs^jkFE2K?RmmqbFHoP?C8`GMc4zOnMubCk;O|l8%5;7Q(~T92 z&TGkFQJ{j>;TSHXIM{YX?M5{Ai`j3we2IfTuhg_Y4&`kCs~_KxXTnaP{Ip1HyS7iS zdN$ngwX4DFi+{DyBdVqf5;74006+&st8asHB<;wg5dZ*3L;wK)U)S2zO49?jfItat z{tQ%sZ-try@xHRHNAAK&oN$9VPz$ANe! zBY(X=EHxbVBm!y>m)GcN7D|XTS8#JH+hUZx6E}43^+({sVA;d4Dr&dn};rC8y`-?9+j&OAK}M2*L0_DgR_Cy>C9F+r*N^BuMiSQPX}J z@UNeJ2J$emmJkKI_Gb-ul)L=I*^;}+&Frl`N#&w-h3`VDRi8whEQNTZv=WCPt(_ zObn|5A6I@<06-5Pc8)^^C0lWUD_nMiUi+pyYr4T0?il%>kUFueN`lN%9OJP^nbXsB z_fEvRILot9<>McQc2|Y%ZuX<4j*pL5+H0F9&*IjT7e0UVFixACJ3F2mIsNs)6w>X2 zVY)KBTB}}W{Cs))c(G0;oOcvJcTx+In{z&X@@~0!^lmcqxnG}$c62t>d<5Lvo7M#Q zLNrc{I;Wy*{X>K}bUgN;!{{(EEim50`_A6k|rHV-$Oipu8+SeaQE30S;g1Ix)tsm4=Sr z&F_;so>cFY335ey$yyLL8sVbrT&OlkK)NgW;WumOtD#WXG-B-Y=q)_l^LCo(A_0}` zhwtJhM_xN`|HQACu=-5Qvv4G9Ng{G3K5Q=;mng(8)5%++eJ$Fmm*REj065n%vhlP! zEq$a^N8{gx_U{vr7H0@fL1db;;EG#$~zbgiHu#pJK4B<1w?M1~3 zlf;H&xhK@kOOG}j!>(F>e|idM>};TL^l@NWmw#jVPaX?|M$UJ8 zR1_u8%F(~oc=3bHawf;y{feR&cOlp~?yQ25Jx>FMlX<6vG3oh~$P<;+l> zh}UkE4ZdChS^MCqmFKJZ>yo!C8IkK8xzXbA*-y2*) zRijMQ1Y1wzU>it(Bt|MxN;;|u0Xah^=cM=%Zzk#70{wSK4pDB&oeL zMoOPl09$>Py*4HT@t$`sT9uL)E&5T%tYIYmFp8TR69 zNTh(e?Z}41f4#9{^32EO*RqbG;l9+C4@_Xa{G3tf&yer;CJxLub=^Sb3niydMw>Xs zVxWlL`5MOk%#*tLXyzJ)bQ*4)A60r(l}}P@uq38`Tphh77fxNbe!6;4;Y9NEv)9ta z%&IkdEv)N%KAve!HAgtJuXIAGj%!r67;jSX^09PhdDCiEf1&2`(?-t^idv$b;E3_! zyk_b=mvS3>Ekz4A_4dY4HE65Z=9)#DX*-2TlUdxvvV_a7hs=J-> z^Vmrdq3uoQ-cQ`3w7Dq9_ZD=SG@%Z`A_IRBckIY=YYFp226ho)c9NM!=9fGD?oH7) zi@G+~4DX_}WQ zpp3W^A^o7`2j=~Rz4ixh0N3(xz#)T7NW0>o2M?E^@9>MqpJ(HBE`vqirM2dmGhK^5iIVWHik3;S43O08q*U9A@@ZJ*M5By_9jw7)RYkqX*CL6Dqc^OvHfjB;gjEe>xZHdx22Fesug3gw(68$O&2c@pcWs*1=!@^PZrY- zMf9M$0-LJV673=Q7w;yutw(t*kT>I8_w3dLT}?N4V7S z?8b`GF;lB@`C0f;!^@<9727-SvicW<)HNEmQMp)IPzsXGwJQ}u7LXuxvyY-k-K}7X z{~l99sSkfL=ej)p8nA$(z_h$dRR6)m-SwmGRa7V{V_-3Ka>I)Q&)((gv0wIF@V1SD zARu6<-aWq**%Abv-nyhE^1So{G~O(ea!O5_!JK7N4(djCpP&|H>o8keBYcD*vDxfIUUph{qVdU~xdE}r$61r*$L?Jsx zPew^tU-g-d0cPV*YlT0E; zj04p}UBScKGH)u=@{Om!XHr~pyYJ@?7ov{cK@QHYut;9VkFl$229A!4Q&I-74Jh1~ zF9aJZY=2!iTEW{eX;#XGJGe}$G}Bj6OwCTcP4Bxz?w-Gb7>8>Z^%WQs*WFsiJxV4( zJl*%RNgKiwoUIl+TewJ{CkMW7?q@vEFx5P!SjP06@EAx^r$$yl(o{a)LN^CB;(?`nYIiRvt4R+Ag7s+N^kE z&}`6=pWAaXJVzppM1orMc{1#xryc6*IK%m#X8%dW_Ej??-EOoxdhzUH0?oDXVzgRu zxATK4ht^^qF`kSjyQ*>cy{a?Q=%1F zqnfOnZS)&4?2?AgprP$vhs1u&LsVheZ-@|kxe8J2LZ-s@6waQ>z*m($HW1XkBm;3` z+V?Bi)2cRT>vk8Ag3Dxp31BQMmvo-g7N+~$HI*Og6ngA#MqK@{dU~2KY_RGra^aj!+`gI-@~ed(JNdZ`caUM2@+? zgHPDfo3@BCUg`YCy-e*(eSByLeE^TZ8%gN-x>MEh2;*nloF4KoCslZ5k}RL8teZTpZ`n|>lGAdG zVGj@EwMMBBc7~VY)egg3*YRnPQJiJzZ5&>TXN3ZMweZ<0 zbVFx{Ubr~`wL5Ryq_Gn!VaC)mS1ylX{?3PTz{OOky8?{ATPFGUOp{0R&5%H=5IkjU zE$L%DpHvIsKt!E4hQM-kv0QGC1t!|Ty6j*#mm*+*TJpBp4^@Gnz|U(F73zEjzV~b^ zptFL%8IKz{$l47LXtC+F%e=R2RxFhymkdu_yhZPbLCT#VN^^i7ZtheQmesk5g(`H; z5cBt6+?NsLPGV}QnIy-O+C2Se9)ci*lhoN3z_Tr*$fy|H-D-WSCK@jOS%=kqDC?0^CXt>mI14j2MRw2KH%v ze(Z}m9#poxtUp0(rW3G@&w&nXB-c9?h33B4d=XBF>BW3AgTk zM9_w|p4P($yDJ?`QNrZ@kdv`h13eAAihWl+r zJu5V=Ue?O}tS?U^jw3PWdgb{-MIShGg2j_fc=uzm9teb@#%sB8;2E1-j-YdTg%UG~ zqi!V_W{6(I?4TzyiX?`x1MyZ=N;Upey>ZJWq&j3sY!v`<#F&Vx+;EOC8FoFr|X!R>yNEPO^SoZhr# zqwDn4U-<%!^BJyA%=iXV;kNU4$r=19^Si%lb?!1dyjj9ZX;A)choGtY7?f~*gSdZ;#@x(B# zkxF1ZSY#}XyJ1s>^Hvg+sfdlCfE^Z?H0qzj?Q@i=>dD6Hq*+B?OF>iN_Jp|Af9{ve zv)OyV=c^IvZu?g&85)6)gz6_gl6WaW1%)TBrPoV;kG1N~+N(ae|25J9ADoda3_Qm9 zCV>_JGCskL8~J#;EIIemj_qrsjxa{k_5*j{O5Lk);j9!<5!t#u>+#C~W4l;2Sooyb z?z8$f@H}f|ymI%Qx7?R+BDOW%LK~b}*MwPinAU!%b}`c++;Qy?c9DDdlOK)q0(8?K zxI2>fC*;Y?it2A<>6x1(@!Ac?oda|X+0Q=@n3>dl9N=hk(=)Guf0hErxd z*k|&)X0HiW4aJE=9;H{KHURF*sYwPSxW3rsJBkS77&X$;Q6{>U|)g8;+lyCM%g z^bbD!B2CQ)doekw44nS0^b)ni-kzNWKca9h9+}CpAn)jiQ5o8;poYv4j|!PMU?7da zh#ew8MJH#C;cr6Zg-TJp~_RvX-%@x;RnmT^{0t=ifF|EvBPX=8o<+S)0DMCXvGeI+C*Lf zABdzGHIwM17T=@RzjXP&9VV39@-2AWL<2N7KU@aKGY)fkgOx3-^UgYdkzje*ky(q znV2e^m^?|)JP*%3DL(QZd6xC}c}Wg63df0~9>WEuv;S*c-%colxf6Nf2cn9B-^O|N zqLs^c6BIgkQD<(oyOw8|KFg!)^5Fl8dU6NTkfD=Qh96?-%rE}TZ4m&AJGO#g|;_0WBO`E}QqUwbW4_glvx5_05Iq zh+r+^E|nYQ7#RwH&K?z}H!BM(F1x;E;rz>b-*2`Gf?QV$*;baYQbk-J<>?t5oL?5x z={}HbRI{P3=S!}|v;`rdcDEpzTQJAT-gPMZxPEkT>Jz{Ng!dZ;Fev{mVzEe=DCo>n zuy`p=jJ@1$w2mmeSCcmKYX30ZVJO|~Mb~`2<+f5R7xlAXEsEejzCQO}s%eWuniA1X z9>$~~g%~Ievr*tfERcp}DG1{bZ~_2U5Kk{xX9H&^ZkUz|73@I)2XR;i)*wv;qful- zype%zNYlZx!FVuVMG8bwPFR|v6e8dq?Bp#4>_SoT|0oYkObJN&|KYVCYcq$8%1q&k*%caWF{70PXzmdrQk&FK>$5FzhDG6aKU=GA1 z9vF(UB;qVTOkP&gv8YG66?w0Njf%Bd7 zU;JyGn|-m@-s`>I&zq8kvXzdKB#dEP4;i%lCP54B_Ur)?DV?T15bB=9$&2LayB_JjP)z=ZyF#vxiaL%t&#*QJk_Z zU2mw#+YnDVi4{cT=$9kGUQnF=K(Q~r ze7lCo5qcUMItXonOW!OO=KwR=cm#OzFW1ZuKA4^KTrAq!%}yC0plK|-Blk@~+~`dq zO4dV7kUS!1CMn1w(7q8P>_jE|*#Z1Zrvgg;41SrTsXdOUKxCa0>}iRgA@vT;ymx}v z{5uyg>1A_il38y44&cx^zaDl^E*&h+8A`!r z=Ic((%T?0-IuAuRGm~CA8ffUG0000R5T&&#H^5~{p@0GaxS;|7#Q(YWUUoXLfp_?@ zbreGAciQZ5L@!#aFN2A08!obmXP$*MCgAZFVQC!ox()C2QX1knm>y)QYQ5f8S0X3G z$brgSY2auMPj@9-gAZbPveZ<4My2g4a)5kFjf=!NTt_ZoTIj47joC}Re*ZauQ3%hD zF96=6Ez#q@(JwKo*yL}%?0w$d6soVVW1cshfQ@W^gh&?!39K3{q`hlsgWYC?@s zbuu;h-3P+a=;Ty#Lus2NS!))a6dM&B)F-*wE;CBV7{BVN3dJ`GRUZC%yPsdYBU9P>W|<=i9JO}F3N)PAn&`BA!=Gzt-z_lY7*}7-E}yv-cd8qL zR?mpmjr(fu*WZE>3;+N@2sn@Nc>_Ia(i>Twn~#GJx4YFfUCp!;9y>zlRLZO`*UxY7V2dRWtHRw zEUV85;LFbR7CIdGtXZbj9cNUFg7go!di*R8G+noc7Px-q7+4@cW%hhci_zSF1I*Kd zA4D~NrA~Np@2A_y`Nl)a< z50p2ZmIYCBlUYxLICPf2LN^eVFQAClc2RgZ^{K`on}5{mik2X%LUGPI)z){@aE8bV;7{VqyXXtKWcY5Ijk$Ok z8huF~cxgIxvsru8%vrEJF<*Vijn}n=Uo$&kb%{t)xJC0plkdTX+FyK`K5W^xaiO+jES9mWEKzy<4&Z(d`^?5`L!$0U<^5Ol*&ulMCceWdyD@FWJ>jEZ zgueD9?=S3^m`gZ83x<2{Um`HEivEVJ=&k=Z5Tvx9i zj9c^?3Y^s)csVyLnVq>=$sS$Bhkg}><>kG(=aC$0r~5bnZFP}e^(2DZ@c7TF<$6|1 zn)+bD+smyN13IocBvn*w+}SQ9L|;T$3XF6-WyG3i2QZV_u6~8={QR-1NGiPi*)9^g z@`YZ!E8^X#YIEypA}ycT7w(>_mlq0ZS0f=1*tmDOj9Rzb|DUm4T(2*8`o82{gW*%~ z9bT_QdUx(KwBLdP@mfKUMsF?bBgci0yK-!p3I<}(CN>Pufjk^e<0Mwg;6n%bkr8=H zy(XaPobs1EQ>v%qaDe>#U|bg=j}>#cU9wIb9#ixKT^(ZzIjQoW7FyED0p2U()x{Ow;UuKHBAN4jo(eo28>A}Xwi;yKy9K)M%Q0!+73-7y;9j^k}>$q zLNA`;c23BkH>*uN6r(%q-}yVrg1)!X zQN@vSYqbr_byImmeSgy%O_{_FLCbTsr?RKKQ-ms1GIQ}}*{aN)c9x=^P1#rtzEt&< zjLan)#u$kM8v{Gt*+8$Z_4MvY`bwd0Bh*`eAl?>nb4hPZq!X3QpLG|Zi*K(-^Ie`} zRGg@!tO$o`xi!km`k39rFaJz5)wyhcHl7po`J3~tp9QaVBAz4rQ(uAnI{u%v=*Gy> zdgR%D_tOo15>>qpnHuIR^L7M7JRtyJ&c)w>prO@q2(tm56W}^d@bFrhWR+-eP7A{+eYyn_UHfhHEr7iE`>Jq%%i8S_({EM8 z9FqPcTq=jy1eC77Jtl}dGQOxXMjBh6}-wWH- z*Z89S74n{Hb)PwMD>j|CAj%x~;8GUCCnV;8jl-af=nC76^;b!8h~=uJPwFr&_vG{8 z!`Mk*Po;iiKjUTMK%X^bxG|xXd*F*H!eZ`VF24jmT5=p(6E{3y(936HaKiUcyM<;K z4d)w~H3cI1a9LZwrbQ%>1!tKWpU)GVBY-Gred7sf6JzIOXV{BLHz(+Umx%#KO& z?SuS2b}W2`g?T<)8~*x4OxzzvvVT*eK4O%AW*SgMJ}PI*G7(PQ_l0^Z2W3k z32;)B{8pi(W<2+#!1BB5Z@*qAeECG=CM6n)9U<@ZJ`(iQmjxv39@RG+yd2H$_3SzlA z2KN#WOjY60E|kfv{b#;_&hjn$7o{u*n$*ni+7{;3tbGPU4033;4fSZ2X}J+{fz!hcvpZi9BjnCdZ%}g&k zyB^P79b4-mDFq(GErJx%b|#TLyHA{^L0U#4$SAP=U){i+M(8IeDM3+fBN=;7S+(&e z9S#zD(GtNTzp>fyOPjHd=Zuj~Hr%4I14DuPo`pdBaqd!)otzmAHI@FgQV}nC#+7z% zvU?YaC)?B!8I7tHHQ8WX7v#(ua)`zh;;P*dWo+_6i?VXMQ>D|cN3DM#C$-VQcQK?z ztz;nAX-fx(6q*Y0hG}+aRMnJyUTto81? zX@;bJq;};o3#{`_XgXZZJLr30om%NoH^kaw43XP5!|xRXg|`^{F_ywgKGyr_{kS)| zb2ss1k4VfUf8)RtOWEyPN7-U!X!<*7{raQet%r0LG-9YX4~jgZkDZvLfq=$aZgM~E z@?1$nK>5M9aqn4&S~~C)usI`K+<^IIjVz`Ql*+k|j7ZmH4lg-_#OfcRkB5J_zY|?~ zJ^a0#Ud+DgU%7(u?|~20?7*$DeR(r95pO2&Ln75NT-Ty?Jg5}jDflp(24>$sryuBs zZVp;NUjY?icJIxyFOPAX*n`f38M+q%6#UZHK6VKNXnM${Bg#SYA01;I)hC`wyy(=~ zmp_&`Vy5)z2lMUD{xRr@Ysbv0P;lyGH@aYxvc^sA<0si)XGsC9c$3QLEy) zM_Rm2%U<6Zxk0aSdLk(wK*tzLvr}g@Q3OOnkGdyVzP_F8tb(b|eaNX(wcj`Yy`u7p zl~k-S$Xb7wkg!=H%14J)>q{-)Z}>!g3WiVVZw&8k^?#$8BR%?QfN80sj}JFqKS(G% zi@!`2CVO0eQ6v;Awf@jS(!5u%LW;aRZvZAl47`M?;(%+N;`e|juOqofp)2pPzTV>$ zLeniHqmI}=SYhvIx;SK=kR`%0UX9qYFcp?KuLqPk^i)BLndf82-&+#xQ^3>->P(9v zkHWZYYMi4UgI|8eR-_r~C#b*tXz)yXS{kbRa=$qH^rW5$EX3cKv&RuAj+&?DUg{_>PLrL!2Wq5P}ZY^sR0f+VbriuwnqWC>t>5&QQ z93?I7GuH;5P|e)bh%#w*ke%v6$HH3G>OdI4Xtb%g@qbK_Cn;TJ&1< z?$A|xD1Qk9gL3J;$v66rE|$|z2IeTik*Rz}rA&n7XO#X3gf);gXWhO8ZlQ<;rMIRQ z{gk?dBea-Moq@t@YxGv0@=3!7MYBO6X8P-|gJeJXVtr>x2*2AnM$R*S2}Tzf^-{YS znGiO@fR06SzxS!L+lGMUg8fQ7i9-Zwg7#X^ofs1fKZ;V~CLmwZe9?X#5X#(VGdlNi zminy>b>mC3*6fLg5R*PCZpQDpz6zz!TUxh*#YL_G`b(X@8iS6vyGCg|J26%SaS zM-kCzS$b72l`N|^U(nKUv;IXpmXly7l0>Xlv-R?#R6QN)U^4e%=!eZUrY4%%7eI{7l_X}!5u zV6z}e$Gv+pTz9?%|YZZnyIZp&V203g*v^xFm8?h6F+y~uewAN*vrK1iyCf( zCspu(y;=xY$FUUX;2$GAFa$5GJvSy;2jQF=WQ3zGaDQ(0a4ul_1lYjLiLpjLt^t7C znDlSo`z5><(}J#xeclQu`An$WO2bzrEv6-wB34lPxV+!d|2?o}r^!*P&zD+{1S`o0 z`yXg7-d+SUZo18`slcKyBiQZ4b~mDw`5{SnApD*f7LTwvS>7=Ds}nN^x8Sl66w`6I z#ZbSv0pS7HboTWyub65-^3B8uCgJb>R(t_UTDrejv_MB9 z9jKW59O<}=lZ{P@7#@+~^%sjF_0Hqk4tcyu#jRpE%0F3#_3VRnJ~Z!!O%;gWw{Iy0 zWc~UP{_kS|VN$+ldzAIO6c#buMlH!oU1EIT$XxUcFEASInh15^*GAj|4j-sMM9I;X zp9(Tja|F#~m`VQE27m@!uw4??J5TiyqEvQ>&B40Jl;CfT3tb(Lhdr=eZ=@c=0FU~B zYgT^I`5EDP|0og;7xG$KTW`l5RsC#pA`Wy=OAxp5vkR4lu`$sbpO+zlxtST@=g0Jp zd~M^oDH{5{KVa6I=5~XmnGxG$7N;kM3OG_938f5jq*PfO+Y_Vy0?R(zW>@P6-OJDY z9aX*y^%TYkO>3!FP98VZ+Ww+t-V4~e=4&VBzJum~wK6F;&rn(o`x?vh!_hnS{vk^$I0YgTwZW#}M)Rf}UDI4+2k7;$U>yUZP+gpoWZ^U~lMROwcsavQ)_V_agH&&7# zN5DjZiP)Rih+||cBXF+l+6g$kLG2+l=iimO+N@mN!8rBq@tL#bKBo<8SHq_e2-d;e zUD6&|6qJP_!NN{>^tJbkUC^GH`)_fdpdY)kquxn|DT5o<6)?iRpG_p7>ePGcp3iD~ zj*=8E^|Mg6OJX|KEC&!%3ocJ?)jt!&DzC zjDa@Q)UCR-vA2@|=Q@GPkJu5_#%f0-8xpON-E+GjiUMJ(&~YrV5DD+p7<~2;N-C6x zS#BAN+g~iB0pw;Vra8u*%qT7=!fHZ7$6eVm&Xb!s(~705RGv}{D0Qvi=h7Q@HG)~L znP)O=6+paOcbMqYu8&4zx_7Y1_CYUke!&wv4g3fnK7pSRsxuf!0uQgh#)iMTgN2$- z@mZJ5s-XOO*j8>3ki?gK&_hD=Ckf{`7q-B)?eF7DhI0XvIvrQdZvO1lUcvoA} z-<#Ejqndy9|MnsQY;2$1Ze+8zWk6B&M5UezBA6)CC|~D?V|CzcH?CRBPe6Y}`Y~1N zP#5=$nU}JOD)NUR_v%m$*~?P4?8n socketFactoryRegister=new HashMap<>(); - private static Map serverSocketFactoryRegister=new HashMap<>(); + private static Map socketTypeRegister=new HashMap<>(); + + static { - socketFactoryRegister.put("TCP", new DefaultSocketFactory()); - serverSocketFactoryRegister.put("TCP", new DefaultServerSocketFactory()); + socketTypeRegister.put("TCP", new SocketType(new DefaultSocketFactory(), new DefaultServerSocketFactory())); + socketTypeRegister.put("UDP", new SocketType(new DefaultDatagramSocketFactory(),new DefaultDatagramServerSocketFactory())); } - public static Map getSocketFactoryRegister() { - return socketFactoryRegister; - } - public static Map getServerSocketFactoryRegister() { - return serverSocketFactoryRegister; - } + + public static Map getSocketTypeRegister() { + return socketTypeRegister; + } private static final long serialVersionUID = 1L; private String type; private String host; @@ -108,6 +108,9 @@ public class MultipurposeSocketAddress implements Serializable{ public int getPort() { return port; } + public String getType() { + return type; + } @Override public String toString() { StringBuilder sb=new StringBuilder(); @@ -125,35 +128,99 @@ public class MultipurposeSocketAddress implements Serializable{ return sb.toString(); } public Socket connectSocket(InetAddress bindip,int bindport,int timeout) throws UnknownHostException, IOException { - Socket s=socketFactoryRegister.get(type).createSocket(); + SocketFactory sf=socketTypeRegister.get(type).getSocketFactory(); + if(sf==null) { + throw new UnsupportedOperationException("Socket Unsupported"); + } + Socket s=sf.createSocket(); s.bind(new InetSocketAddress(bindip, bindport)); s.connect(new InetSocketAddress(host, port),timeout); return s; } public Socket connectSocket(InetAddress bindip,int bindport) throws UnknownHostException, IOException { - Socket s=socketFactoryRegister.get(type).createSocket(); + SocketFactory sf=socketTypeRegister.get(type).getSocketFactory(); + if(sf==null) { + throw new UnsupportedOperationException("Socket Unsupported"); + } + Socket s=sf.createSocket(); s.bind(new InetSocketAddress(bindip, bindport)); s.connect(new InetSocketAddress(host, port)); return s; } public Socket connectSocket() throws UnknownHostException, IOException { - Socket s=socketFactoryRegister.get(type).createSocket(); + SocketFactory sf=socketTypeRegister.get(type).getSocketFactory(); + if(sf==null) { + throw new UnsupportedOperationException("Socket Unsupported"); + } + Socket s=sf.createSocket(); s.connect(new InetSocketAddress(host, port)); return s; } public Socket connectSocket(int timeout) throws UnknownHostException, IOException { - Socket s=socketFactoryRegister.get(type).createSocket(); + SocketFactory sf=socketTypeRegister.get(type).getSocketFactory(); + if(sf==null) { + throw new UnsupportedOperationException("Socket Unsupported"); + } + Socket s=sf.createSocket(); s.connect(new InetSocketAddress(host, port),timeout); return s; } public ServerSocket listenServerSocket() throws UnknownHostException, IOException { - ServerSocket sk=serverSocketFactoryRegister.get(type).createServerSocket(port, 50, InetAddress.getByName(host)); + ServerSocketFactory srf=socketTypeRegister.get(type).getServerSocketFactory(); + if(srf==null) { + throw new UnsupportedOperationException("ServerSocket Unsupported"); + } + ServerSocket sk=srf.createServerSocket(port, 50, InetAddress.getByName(host)); return sk; } public ServerSocket listenServerSocket(int backlog) throws UnknownHostException, IOException { - ServerSocket sk=serverSocketFactoryRegister.get(type).createServerSocket(port, backlog, InetAddress.getByName(host)); + ServerSocketFactory srf=socketTypeRegister.get(type).getServerSocketFactory(); + if(srf==null) { + throw new UnsupportedOperationException("ServerSocket Unsupported"); + } + ServerSocket sk=srf.createServerSocket(port, backlog, InetAddress.getByName(host)); + return sk; + } + public boolean isStream() { + return checkIsStream(type); + } + public static boolean checkIsStream(String type2) { + return socketTypeRegister.get(type2).isStream(); + } + public DatagramSocket connectDatagramSocket(InetAddress bindip,int bindport) throws IOException { + DatagramSocketFactory dgs=socketTypeRegister.get(type).getDatagramSocketFactory(); + if(dgs==null) { + throw new UnsupportedOperationException("DatagramSocket Unsupported"); + } + DatagramSocket dgd=dgs.createSocket(); + dgd.bind(new InetSocketAddress(bindip, bindport)); + dgd.connect(new InetSocketAddress(host, port)); + return dgd; + } + public DatagramSocket connectDatagramSocket() throws IOException { + DatagramSocketFactory dgs=socketTypeRegister.get(type).getDatagramSocketFactory(); + if(dgs==null) { + throw new UnsupportedOperationException("DatagramSocket Unsupported"); + } + DatagramSocket dgd=dgs.createSocket(); + dgd.connect(new InetSocketAddress(host, port)); + return dgd; + } + public DatagramServerSocket listenDatagramServerSocket() throws UnknownHostException, IOException { + DatagramServerSocketFactory srf=socketTypeRegister.get(type).getDatagramServerSocketFactory(); + if(srf==null) { + throw new UnsupportedOperationException("DatagramServerSocket Unsupported"); + } + DatagramServerSocket sk=srf.createDatagramServerSocket(port, 50, InetAddress.getByName(host)); + return sk; + } + public DatagramServerSocket listenDatagramServerSocket(int backlog) throws UnknownHostException, IOException { + DatagramServerSocketFactory srf=socketTypeRegister.get(type).getDatagramServerSocketFactory(); + if(srf==null) { + throw new UnsupportedOperationException("DatagramServerSocket Unsupported"); + } + DatagramServerSocket sk=srf.createDatagramServerSocket(port, backlog, InetAddress.getByName(host)); return sk; } - } diff --git a/src/org/kne/cloud/network/mport/NetworkService.java b/src/org/kne/cloud/network/NetworkService.java similarity index 84% rename from src/org/kne/cloud/network/mport/NetworkService.java rename to src/org/kne/cloud/network/NetworkService.java index 4ffd5b1..9e4a6b5 100644 --- a/src/org/kne/cloud/network/mport/NetworkService.java +++ b/src/org/kne/cloud/network/NetworkService.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; public interface NetworkService { public void listen(MultipurposeSocketAddress msa); diff --git a/src/org/kne/cloud/network/mport/PortMultiUse.java b/src/org/kne/cloud/network/PortMultiUse.java similarity index 94% rename from src/org/kne/cloud/network/mport/PortMultiUse.java rename to src/org/kne/cloud/network/PortMultiUse.java index 56d292a..5f327da 100644 --- a/src/org/kne/cloud/network/mport/PortMultiUse.java +++ b/src/org/kne/cloud/network/PortMultiUse.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.File; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/PortRelay.java b/src/org/kne/cloud/network/PortRelay.java similarity index 78% rename from src/org/kne/cloud/network/mport/PortRelay.java rename to src/org/kne/cloud/network/PortRelay.java index 5d2df50..78123b0 100644 --- a/src/org/kne/cloud/network/mport/PortRelay.java +++ b/src/org/kne/cloud/network/PortRelay.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.DataOutputStream; @@ -12,12 +12,12 @@ import java.util.*; public class PortRelay { private int port; private HostPortMap services; - private TCPListener ssc; + private SocketListener ssc; public PortRelay(int port, HostPortMap services) throws IOException { this.services = services; this.port = port; - MultipurposeSocketAddress.getServerSocketFactoryRegister().put("ProtocolDetectorServerSocket", new ProtocolDetectorServerSocketFactory()); - ssc=new TCPListener(new MultipurposeSocketAddress("ProtocolDetectorServerSocket", "0.0.0.0", port)); + MultipurposeSocketAddress.getSocketTypeRegister().put("ProtocolDetectorServerSocket",new SocketType(null, new ProtocolDetectorServerSocketFactory()) ); + ssc=new SocketListener(new MultipurposeSocketAddress("ProtocolDetectorServerSocket", "0.0.0.0", port)); } public void start() throws IOException { ssc.setCon((s)->{ @@ -55,7 +55,6 @@ public class PortRelay { } }); - ssc.open(); System.out.println("已打开端口:" + port); } diff --git a/src/org/kne/cloud/network/mport/Protocol.java b/src/org/kne/cloud/network/Protocol.java similarity index 87% rename from src/org/kne/cloud/network/mport/Protocol.java rename to src/org/kne/cloud/network/Protocol.java index 09098c0..198c136 100644 --- a/src/org/kne/cloud/network/mport/Protocol.java +++ b/src/org/kne/cloud/network/Protocol.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.util.Objects; diff --git a/src/org/kne/cloud/network/mport/ProtocolDetector.java b/src/org/kne/cloud/network/ProtocolDetector.java similarity index 94% rename from src/org/kne/cloud/network/mport/ProtocolDetector.java rename to src/org/kne/cloud/network/ProtocolDetector.java index e369d66..d076551 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetector.java +++ b/src/org/kne/cloud/network/ProtocolDetector.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.InputStream; diff --git a/src/org/kne/cloud/network/mport/ProtocolDetectorItem.java b/src/org/kne/cloud/network/ProtocolDetectorItem.java similarity index 76% rename from src/org/kne/cloud/network/mport/ProtocolDetectorItem.java rename to src/org/kne/cloud/network/ProtocolDetectorItem.java index 5c3ca7b..d704f9b 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetectorItem.java +++ b/src/org/kne/cloud/network/ProtocolDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.InputStream; import java.util.function.Function; diff --git a/src/org/kne/cloud/network/mport/ProtocolDetectorServerSocket.java b/src/org/kne/cloud/network/ProtocolDetectorServerSocket.java similarity index 94% rename from src/org/kne/cloud/network/mport/ProtocolDetectorServerSocket.java rename to src/org/kne/cloud/network/ProtocolDetectorServerSocket.java index 659f91a..4e0b766 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetectorServerSocket.java +++ b/src/org/kne/cloud/network/ProtocolDetectorServerSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.InputStream; diff --git a/src/org/kne/cloud/network/mport/ProtocolDetectorServerSocketFactory.java b/src/org/kne/cloud/network/ProtocolDetectorServerSocketFactory.java similarity index 82% rename from src/org/kne/cloud/network/mport/ProtocolDetectorServerSocketFactory.java rename to src/org/kne/cloud/network/ProtocolDetectorServerSocketFactory.java index fe5b9e4..62571e9 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetectorServerSocketFactory.java +++ b/src/org/kne/cloud/network/ProtocolDetectorServerSocketFactory.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.net.InetAddress; @@ -17,9 +17,13 @@ public class ProtocolDetectorServerSocketFactory extends DefaultServerSocketFact this.pd = pd; } + @Override + public ServerSocket createServerSocket() throws IOException { + return new ProtocolDetectorServerSocket(super.createServerSocket()); + } + @Override public ServerSocket createServerSocket(int port) throws IOException { - // TODO 自动生成的方法存根 return new ProtocolDetectorServerSocket( super.createServerSocket(port),pd); } @@ -33,13 +37,11 @@ public class ProtocolDetectorServerSocketFactory extends DefaultServerSocketFact @Override public ServerSocket createServerSocket(int port, int backlog) throws IOException { - // TODO 自动生成的方法存根 return new ProtocolDetectorServerSocket(super.createServerSocket(port, backlog),pd); } @Override public ServerSocket createServerSocket(int port, int backlog, InetAddress ifAddress) throws IOException { - // TODO 自动生成的方法存根 return new ProtocolDetectorServerSocket(super.createServerSocket(port, backlog, ifAddress),pd); } diff --git a/src/org/kne/cloud/network/mport/ProtocolDetectorSocket.java b/src/org/kne/cloud/network/ProtocolDetectorSocket.java similarity index 90% rename from src/org/kne/cloud/network/mport/ProtocolDetectorSocket.java rename to src/org/kne/cloud/network/ProtocolDetectorSocket.java index 858e368..48ccf30 100644 --- a/src/org/kne/cloud/network/mport/ProtocolDetectorSocket.java +++ b/src/org/kne/cloud/network/ProtocolDetectorSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.BufferedInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/ProtocolStack.java b/src/org/kne/cloud/network/ProtocolStack.java similarity index 65% rename from src/org/kne/cloud/network/mport/ProtocolStack.java rename to src/org/kne/cloud/network/ProtocolStack.java index 13d793a..58b1935 100644 --- a/src/org/kne/cloud/network/mport/ProtocolStack.java +++ b/src/org/kne/cloud/network/ProtocolStack.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.util.Stack; diff --git a/src/org/kne/cloud/network/mport/ProxyProfileEntry.java b/src/org/kne/cloud/network/Proxy.java similarity index 65% rename from src/org/kne/cloud/network/mport/ProxyProfileEntry.java rename to src/org/kne/cloud/network/Proxy.java index 134871a..e0e09fb 100644 --- a/src/org/kne/cloud/network/mport/ProxyProfileEntry.java +++ b/src/org/kne/cloud/network/Proxy.java @@ -1,5 +1,7 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; +import java.io.Closeable; +import java.net.InetSocketAddress; import java.net.Socket; import java.util.HashMap; import java.util.Map; @@ -8,9 +10,13 @@ import java.util.function.Consumer; import javax.net.ServerSocketFactory; import javax.net.SocketFactory; -public class ProxyProfileEntry { +public abstract class Proxy implements Closeable{ private static Map register=new HashMap<>(); public static Map getRegister() { return register; } + + + + } diff --git a/src/org/kne/cloud/network/mport/ProxyProfileAnalyser.java b/src/org/kne/cloud/network/ProxyProfileAnalyser.java similarity index 100% rename from src/org/kne/cloud/network/mport/ProxyProfileAnalyser.java rename to src/org/kne/cloud/network/ProxyProfileAnalyser.java diff --git a/src/org/kne/cloud/network/mport/ProxyProfileExecutor.java b/src/org/kne/cloud/network/ProxyProfileExecutor.java similarity index 100% rename from src/org/kne/cloud/network/mport/ProxyProfileExecutor.java rename to src/org/kne/cloud/network/ProxyProfileExecutor.java diff --git a/src/org/kne/cloud/network/mport/RDPDetectorItem.java b/src/org/kne/cloud/network/RDPDetectorItem.java similarity index 86% rename from src/org/kne/cloud/network/mport/RDPDetectorItem.java rename to src/org/kne/cloud/network/RDPDetectorItem.java index 5d0d38c..dbaa497 100644 --- a/src/org/kne/cloud/network/mport/RDPDetectorItem.java +++ b/src/org/kne/cloud/network/RDPDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/SSHDetectorItem.java b/src/org/kne/cloud/network/SSHDetectorItem.java similarity index 87% rename from src/org/kne/cloud/network/mport/SSHDetectorItem.java rename to src/org/kne/cloud/network/SSHDetectorItem.java index 280c071..ebe068c 100644 --- a/src/org/kne/cloud/network/mport/SSHDetectorItem.java +++ b/src/org/kne/cloud/network/SSHDetectorItem.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.DataInputStream; import java.io.IOException; diff --git a/src/org/kne/cloud/network/mport/ServiceElement.java b/src/org/kne/cloud/network/ServiceElement.java similarity index 91% rename from src/org/kne/cloud/network/mport/ServiceElement.java rename to src/org/kne/cloud/network/ServiceElement.java index 3bb5cff..6fd86b7 100644 --- a/src/org/kne/cloud/network/mport/ServiceElement.java +++ b/src/org/kne/cloud/network/ServiceElement.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.net.UnknownHostException; import java.util.HashMap; diff --git a/src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java b/src/org/kne/cloud/network/ServiceToSocketProxy.java similarity index 53% rename from src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java rename to src/org/kne/cloud/network/ServiceToSocketProxy.java index acf96d5..6299a34 100644 --- a/src/org/kne/cloud/network/mport/ServiceToSocketProxyProfileEntry.java +++ b/src/org/kne/cloud/network/ServiceToSocketProxy.java @@ -1,18 +1,24 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; +import java.io.IOException; import java.util.function.Consumer; -public class ServiceToSocketProxyProfileEntry extends ProxyProfileEntry{ +public class ServiceToSocketProxy extends Proxy{ public NetworkService getSrc() { return src; } public MultipurposeSocketAddress getDes() { return des; } - public ServiceToSocketProxyProfileEntry(String l, String r) { + public ServiceToSocketProxy(String l, String r) { src=getRegister().get(l.substring(1, l.length()-1)); des=new MultipurposeSocketAddress(r); } private NetworkService src; private MultipurposeSocketAddress des; + @Override + public void close() throws IOException { + // TODO 自动生成的方法存根 + + } } diff --git a/src/org/kne/cloud/network/mport/SocketBridge.java b/src/org/kne/cloud/network/SocketBridge.java similarity index 55% rename from src/org/kne/cloud/network/mport/SocketBridge.java rename to src/org/kne/cloud/network/SocketBridge.java index e4cedc2..27c0469 100644 --- a/src/org/kne/cloud/network/mport/SocketBridge.java +++ b/src/org/kne/cloud/network/SocketBridge.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; @@ -6,33 +6,37 @@ import java.net.*; import org.kne.io.Task; public class SocketBridge extends Task{ - private Socket a; - private Socket b; - public StreamBridge getSab() { - return sab; + protected Socket a; + protected Socket b; + protected StreamBridge bridgeAB; + protected StreamBridge bridgeBA; + public StreamBridge getBridgeAB() { + return bridgeAB; } - public StreamBridge getSba() { - return sba; + public StreamBridge getBridgeBA() { + return bridgeBA; } - 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()); + createStreamBridge(); + } + protected void createStreamBridge() throws IOException { + bridgeAB = new StreamBridge(a.getInputStream(), b.getOutputStream()); + bridgeBA = new StreamBridge(b.getInputStream(), a.getOutputStream()); } @Override protected void runTask() { try { - sab.runAtNewThread("SocketBridge A->B thread"); - sba.runAtNewThread("SocketBridge B->A thread"); - sab.waitfortask(); - sba.waitfortask(); + bridgeAB.runAtNewThread("SocketBridge A->B thread"); + bridgeBA.runAtNewThread("SocketBridge B->A thread"); + bridgeAB.waitfortask(); + bridgeBA.waitfortask(); } catch (Exception e) { e.printStackTrace(); }finally { + //new Exception().printStackTrace(); try { a.close(); } catch (IOException e) { @@ -53,7 +57,7 @@ public class SocketBridge extends Task{ public Socket getB() { return b; } - protected OutputStream getBOUT() throws IOException { + /*protected OutputStream getBOUT() throws IOException { return b.getOutputStream(); } protected InputStream getBIN() throws IOException { @@ -64,6 +68,6 @@ public class SocketBridge extends Task{ } protected InputStream getAIN() throws IOException { return a.getInputStream(); - } + }*/ } diff --git a/src/org/kne/cloud/network/mport/TCPListener.java b/src/org/kne/cloud/network/SocketListener.java similarity index 52% rename from src/org/kne/cloud/network/mport/TCPListener.java rename to src/org/kne/cloud/network/SocketListener.java index 4a9da4d..b222819 100644 --- a/src/org/kne/cloud/network/mport/TCPListener.java +++ b/src/org/kne/cloud/network/SocketListener.java @@ -1,22 +1,25 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; +import java.io.Closeable; import java.io.IOException; import java.net.InetAddress; import java.net.ServerSocket; import java.net.Socket; +import java.net.SocketException; +import java.net.UnknownHostException; import java.util.function.Consumer; import javax.net.ServerSocketFactory; -public class TCPListener { - private ServerSocket serverSocket; +public class SocketListener implements Closeable,AutoCloseable{ + protected ServerSocket serverSocket; public ServerSocket getServerSocket() { return serverSocket; } - private volatile boolean flag=false; + private volatile boolean flag=true; - private Consumercon; + private volatile Consumercon; private Runnable r=new Runnable() { @Override public void run() { @@ -24,8 +27,26 @@ public class TCPListener { try { Socket soce=serverSocket.accept(); ThreadTool.makeVThreadIfSupport("端口监听线程",()->{ + if(con!=null) { + try { con.accept(soce); + }catch(Exception e) { + e.printStackTrace(); + try { + soce.close(); + } catch (IOException er) { + er.printStackTrace(); + } + } + }else { + try { + soce.close(); + } catch (IOException e) { + e.printStackTrace(); + } + } }).start(); + }catch(SocketException e) { } catch (IOException e) { e.printStackTrace(); } @@ -34,14 +55,17 @@ public class TCPListener { }; private MultipurposeSocketAddress multipurposeSocketAddress; - public TCPListener(MultipurposeSocketAddress multipurposeSocketAddress) { + public SocketListener(MultipurposeSocketAddress multipurposeSocketAddress) throws IOException { this.multipurposeSocketAddress=multipurposeSocketAddress; + open(); } - - public void open() throws IOException { - flag=true; + public SocketListener(ServerSocket tserverSocket) throws IOException { + this.serverSocket=tserverSocket; + open(); + } + protected void open() throws UnknownHostException, IOException { + if(serverSocket==null) serverSocket=multipurposeSocketAddress.listenServerSocket(); - //servers=new ServerSocket(port); new Thread(r).start(); } public void close() { diff --git a/src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java b/src/org/kne/cloud/network/SocketToServiceProxy.java similarity index 55% rename from src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java rename to src/org/kne/cloud/network/SocketToServiceProxy.java index a0cc947..dc407ac 100644 --- a/src/org/kne/cloud/network/mport/SocketToServiceProxyProfileEntry.java +++ b/src/org/kne/cloud/network/SocketToServiceProxy.java @@ -1,10 +1,11 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; +import java.io.IOException; import java.net.Socket; import java.util.function.Consumer; -public class SocketToServiceProxyProfileEntry extends ProxyProfileEntry { - public SocketToServiceProxyProfileEntry(String l, String r) { +public class SocketToServiceProxy extends Proxy { + public SocketToServiceProxy(String l, String r) { src=new MultipurposeSocketAddress(l); des=getRegister().get(r.substring(1, r.length()-1)) ; } @@ -16,4 +17,9 @@ public class SocketToServiceProxyProfileEntry extends ProxyProfileEntry { } private MultipurposeSocketAddress src; private NetworkService des; + @Override + public void close() throws IOException { + // TODO 自动生成的方法存根 + + } } diff --git a/src/org/kne/cloud/network/SocketToSocketProxy.java b/src/org/kne/cloud/network/SocketToSocketProxy.java new file mode 100644 index 0000000..836020b --- /dev/null +++ b/src/org/kne/cloud/network/SocketToSocketProxy.java @@ -0,0 +1,130 @@ +package org.kne.cloud.network; + +import java.io.IOException; +import java.net.InetAddress; +import java.net.InetSocketAddress; +import java.net.Socket; +import java.net.UnknownHostException; +import java.util.Map; + +public class SocketToSocketProxy extends Proxy { + private SocketListener sl; + private DatagramSocketListener dsl; + private MultipurposeSocketAddress listen, cbind, defaultConnect; + private Map detectedConnect; + + public SocketToSocketProxy(MultipurposeSocketAddress listen, MultipurposeSocketAddress connect) throws IOException { + this(listen, new MultipurposeSocketAddress(new InetSocketAddress(0)), connect); + } + + public SocketToSocketProxy(String l, String r) throws IOException { + this(new MultipurposeSocketAddress(l), new MultipurposeSocketAddress(r)); + } + + public SocketToSocketProxy(MultipurposeSocketAddress listen, MultipurposeSocketAddress cbind, + MultipurposeSocketAddress defaultConnect, Map detectedConnect) + throws IOException { + this.listen = listen; + this.cbind = cbind; + this.defaultConnect = defaultConnect; + this.detectedConnect = detectedConnect; + open(); + } + + public SocketToSocketProxy(MultipurposeSocketAddress listen, MultipurposeSocketAddress cbind, + MultipurposeSocketAddress defaultConnect) throws IOException { + this(listen, cbind, defaultConnect,null); + } + + private void open() throws IOException { + if (listen.isStream()) { + if (detectedConnect == null || detectedConnect.isEmpty()) { + sl = new SocketListener(listen); + } else { + sl = new SocketListener(new ProtocolDetectorServerSocket(listen.listenServerSocket())); + } + sl.setCon((sox) -> { + Socket sk = null; + try { + if (sox instanceof ProtocolDetectorSocket) { + ProtocolStack ps = ((ProtocolDetectorSocket) sox).getProtocolStack(); + if (!ps.isEmpty()) { + MultipurposeSocketAddress pmsa = detectedConnect.get(ps.pop().getName()); + if(pmsa!=null) { + sk = pmsa.connectSocket(InetAddress.getByName(cbind.getHost()), cbind.getPort()); + }else { + sk = defaultConnect.connectSocket(InetAddress.getByName(cbind.getHost()), cbind.getPort()); + + } + } else { + sk = defaultConnect.connectSocket(InetAddress.getByName(cbind.getHost()), cbind.getPort()); + } + } else { + sk = defaultConnect.connectSocket(InetAddress.getByName(cbind.getHost()), cbind.getPort()); + } + runBridge(sox,sk); + } catch (IOException e) { + e.printStackTrace(); + } finally { + if (sk != null) + try { + sk.close(); + } catch (IOException e) { + e.printStackTrace(); + } + try { + sox.close(); + } catch (IOException e) { + e.printStackTrace(); + } + } + }); + } else { + dsl = new DatagramSocketListener(listen); + dsl.setCon((dox) -> { + + }); + } + } + + protected void runBridge(Socket sk, Socket sox) throws IOException { + SocketBridge sb = new SocketBridge(sk, sox); + sb.run(); + } + + public MultipurposeSocketAddress getListen() { + return listen; + } + + public MultipurposeSocketAddress getCbind() { + return cbind; + } + + public MultipurposeSocketAddress getDefaultConnect() { + return defaultConnect; + } + + @Override + public void close() throws IOException { + if (dsl != null) + dsl.close(); + if (sl != null) + sl.close(); + } + + public InetAddress getListenAddress() { + if (sl != null) { + return sl.getServerSocket().getInetAddress(); + } else { + return dsl.getDatagramServerSocket().getLocalAddress(); + } + } + + public int getListenPort() { + if (sl != null) { + return sl.getServerSocket().getLocalPort(); + } else { + return dsl.getDatagramServerSocket().getLocalPort(); + } + } +} diff --git a/src/org/kne/cloud/network/mport/StreamBridge.java b/src/org/kne/cloud/network/StreamBridge.java similarity index 71% rename from src/org/kne/cloud/network/mport/StreamBridge.java rename to src/org/kne/cloud/network/StreamBridge.java index 73f1c7c..ea5d97d 100644 --- a/src/org/kne/cloud/network/mport/StreamBridge.java +++ b/src/org/kne/cloud/network/StreamBridge.java @@ -1,15 +1,16 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.*; import java.util.Objects; +import java.util.UUID; import org.kne.io.Task; public class StreamBridge extends Task{ - private InputStream in; - private OutputStream out; - private int blocksize=65535; - private long delay=0; + protected InputStream in; + protected OutputStream out; + protected int blocksize=65535; + protected long delay=0; public StreamBridge(InputStream in, OutputStream out) { super(); Objects.requireNonNull(in); @@ -20,12 +21,16 @@ public class StreamBridge extends Task{ @Override public void runTask() { + //FileOutputStream fos=null; try { + //fos=new FileOutputStream(UUID.randomUUID()+".txt"); int v=-1; byte[]b=new byte[blocksize]; while ((v=in.read(b))!=-1) { + // fos.write(b, 0, v); out.write(b,0,v); out.flush(); + if(delay>0) try { Thread.sleep(delay); } catch (InterruptedException e) { @@ -36,13 +41,21 @@ public class StreamBridge extends Task{ // TODO 自动生成的 catch 块 e.printStackTrace(); }finally { + /*if(fos!=null) + try { + fos.close(); + } catch (IOException e1) { + e1.printStackTrace(); + }*/ try { + if(in!=null) in.close(); } catch (IOException e) { // TODO 自动生成的 catch 块 e.printStackTrace(); } try { + if(in!=null) out.close(); } catch (IOException e) { // TODO 自动生成的 catch 块 diff --git a/src/org/kne/cloud/network/mport/ThreadTool.java b/src/org/kne/cloud/network/ThreadTool.java similarity index 72% rename from src/org/kne/cloud/network/mport/ThreadTool.java rename to src/org/kne/cloud/network/ThreadTool.java index 2982542..62238e1 100644 --- a/src/org/kne/cloud/network/mport/ThreadTool.java +++ b/src/org/kne/cloud/network/ThreadTool.java @@ -1,10 +1,10 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.lang.reflect.Method; public class ThreadTool { public static boolean first=true; - public static boolean forceD=true; + public static boolean forceD=false; public static Thread makeVThreadIfSupport(String name,Runnable r) { if(forceD) return new Thread(r, name); @@ -21,7 +21,7 @@ public class ThreadTool { }catch(Throwable e) { //e.printStackTrace(); if(first) { - System.out.println("请使用java19以上版本并开启--enable-preview选项以提高性能!"); + System.out.println("提示:请使用java21以上版本以提高本软件的数据转发性能!"); first=false; } return new Thread(r, name); @@ -46,7 +46,7 @@ public class ThreadTool { }catch(Throwable e) { //e.printStackTrace(); if(first) { - System.out.println("请使用java19以上版本以提高性能!"); + System.out.println("提示:请使用java21以上版本以提高本软件的数据转发性能!"); first=false; } Thread rt=new Thread(r, name); @@ -54,4 +54,14 @@ public class ThreadTool { return rt; } } + + public static Thread makePThreadIfSupport(String name, Runnable r) { + Thread rt=new Thread(r, name); + return rt; + } + public static Thread makePDaemonThreadIfSupport(String name, Runnable r) { + Thread rt=new Thread(r, name); + rt.setDaemon(true); + return rt; + } } diff --git a/src/org/kne/cloud/network/mport/VirtualServerSocket.java b/src/org/kne/cloud/network/VirtualServerSocket.java similarity index 82% rename from src/org/kne/cloud/network/mport/VirtualServerSocket.java rename to src/org/kne/cloud/network/VirtualServerSocket.java index 0fab1ce..80137ed 100644 --- a/src/org/kne/cloud/network/mport/VirtualServerSocket.java +++ b/src/org/kne/cloud/network/VirtualServerSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.lang.reflect.Constructor; @@ -13,8 +13,11 @@ import javax.sql.rowset.RowSetMetaDataImpl; public abstract class VirtualServerSocket extends ServerSocket { - public VirtualServerSocket(SocketImpl si) throws IOException { + private VirtualSocketImpl virtualImpl; + + public VirtualServerSocket(VirtualSocketImpl si) throws IOException { super(si); + this.virtualImpl=si; /*try { Class c=ServerSocket.class; Field con =c.getDeclaredField("impl"); @@ -45,6 +48,10 @@ public abstract class VirtualServerSocket extends ServerSocket { }*/ } + public VirtualSocketImpl getVirtualImpl() { + return virtualImpl; + } + @Override public abstract VirtualSocket accept() throws IOException ; diff --git a/src/org/kne/cloud/network/mport/VirtualSocket.java b/src/org/kne/cloud/network/VirtualSocket.java similarity index 68% rename from src/org/kne/cloud/network/mport/VirtualSocket.java rename to src/org/kne/cloud/network/VirtualSocket.java index 25ef25c..1665fdc 100644 --- a/src/org/kne/cloud/network/mport/VirtualSocket.java +++ b/src/org/kne/cloud/network/VirtualSocket.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.FilterInputStream; import java.io.IOException; @@ -13,21 +13,25 @@ import java.nio.channels.SocketChannel; public abstract class VirtualSocket extends Socket { - private VirtualSocketImpl si; + private VirtualSocketImpl virtualImpl; + + public VirtualSocketImpl getVirtualImpl() { + return virtualImpl; + } public VirtualSocket(VirtualSocketImpl si) throws SocketException { super(si); - this.si=si; + this.virtualImpl=si; } @Override public InputStream getInputStream() throws IOException { - return si.getInputStream(); + return virtualImpl.getInputStream(); } @Override public OutputStream getOutputStream() throws IOException { - return si.getOutputStream(); + return virtualImpl.getOutputStream(); } diff --git a/src/org/kne/cloud/network/mport/VirtualSocketImpl.java b/src/org/kne/cloud/network/VirtualSocketImpl.java similarity index 88% rename from src/org/kne/cloud/network/mport/VirtualSocketImpl.java rename to src/org/kne/cloud/network/VirtualSocketImpl.java index 44a0097..9160906 100644 --- a/src/org/kne/cloud/network/mport/VirtualSocketImpl.java +++ b/src/org/kne/cloud/network/VirtualSocketImpl.java @@ -1,4 +1,4 @@ -package org.kne.cloud.network.mport; +package org.kne.cloud.network; import java.io.IOException; import java.io.InputStream; diff --git a/src/org/kne/cloud/network/klalb/ACKTPacket.java b/src/org/kne/cloud/network/klalb/ACKTPacket.java index b920bc8..d2d1ad3 100644 --- a/src/org/kne/cloud/network/klalb/ACKTPacket.java +++ b/src/org/kne/cloud/network/klalb/ACKTPacket.java @@ -1,19 +1,25 @@ package org.kne.cloud.network.klalb; import java.io.DataInput; +import java.io.DataInputStream; import java.io.DataOutput; +import java.io.DataOutputStream; import java.io.IOException; +import java.util.concurrent.atomic.AtomicInteger; -public class ACKTPacket extends KLALBPacket { +public class ACKTPacket extends KLALBPacket implements PortPacket { private int sport,dport; private long number; private boolean avaliable; - public ACKTPacket(int sport,int dport,long number,boolean avaliable) { + private int sendcount; + + public ACKTPacket(int sport,int dport,long number,boolean avaliable,int sendcount) { super(ACKT); this.sport=sport; this.dport=dport; this.number=number; this.avaliable=avaliable; + this.sendcount=sendcount; } public ACKTPacket() { @@ -31,7 +37,7 @@ public class ACKTPacket extends KLALBPacket { @Override public long getLength() { - return super.getLength()+17; + return super.getLength()+18; } public int getDport() { @@ -45,22 +51,68 @@ public class ACKTPacket extends KLALBPacket { public boolean isAvaliable() { return avaliable; } + + public int getSendcount() { + return sendcount; + } + + private final byte[] writeBuffer = new byte[18]; @Override protected void writeToStream(DataOutput dto) throws IOException { super.writeToStream(dto); - dto.writeInt(sport); + /*dto.writeInt(sport); dto.writeInt(dport); dto.writeLong(number); - dto.writeBoolean(avaliable); - } + dto.writeByte(sendcount); + dto.writeBoolean(avaliable);*/ + writeBuffer[0] = (byte)(sport >>> 24); + writeBuffer[1] = (byte)(sport >>> 16); + writeBuffer[2] = (byte)(sport >>> 8); + writeBuffer[3] = (byte)(sport >>> 0); + + writeBuffer[4] = (byte)(dport >>> 24); + writeBuffer[5] = (byte)(dport >>> 16); + writeBuffer[6] = (byte)(dport >>> 8); + writeBuffer[7] = (byte)(dport >>> 0); + writeBuffer[8] = (byte)(number >>> 56); + writeBuffer[9] = (byte)(number >>> 48); + writeBuffer[10] = (byte)(number >>> 40); + writeBuffer[11] = (byte)(number >>> 32); + writeBuffer[12] = (byte)(number >>> 24); + writeBuffer[13] = (byte)(number >>> 16); + writeBuffer[14] = (byte)(number >>> 8); + writeBuffer[15] = (byte)(number >>> 0); + + writeBuffer[16]=(byte) getSendRecord().size(); + + writeBuffer[17]=(byte) (avaliable ? 1 : 0); + dto.write(writeBuffer); + } + private final byte[] readBuffer = new byte[18]; @Override protected void readFromStream(DataInput din) throws IOException { super.readFromStream(din); - sport=din.readInt(); + /*sport=din.readInt(); dport=din.readInt(); number=din.readLong(); - avaliable=din.readBoolean(); + sendcount=din.readUnsignedByte(); + avaliable=din.readBoolean();*/ + + din.readFully(readBuffer); + sport=((readBuffer[0] << 24) + (readBuffer[1] << 16) + (readBuffer[2] << 8) + (readBuffer[3] << 0)); + dport=((readBuffer[4] << 24) + (readBuffer[5] << 16) + (readBuffer[6] << 8) + (readBuffer[7] << 0)); + number=(((long)readBuffer[8] << 56) + + ((long)(readBuffer[9] & 255) << 48) + + ((long)(readBuffer[10] & 255) << 40) + + ((long)(readBuffer[11] & 255) << 32) + + ((long)(readBuffer[12] & 255) << 24) + + ((readBuffer[13] & 255) << 16) + + ((readBuffer[14] & 255) << 8) + + ((readBuffer[15] & 255) << 0)); + sendcount=readBuffer[16]&0xff; + avaliable=(readBuffer[17] != 0); } -} + +} \ No newline at end of file diff --git a/src/org/kne/cloud/network/klalb/LINESPacket.java b/src/org/kne/cloud/network/klalb/ADDLINESPacket.java similarity index 81% rename from src/org/kne/cloud/network/klalb/LINESPacket.java rename to src/org/kne/cloud/network/klalb/ADDLINESPacket.java index 8e586fe..d941d12 100644 --- a/src/org/kne/cloud/network/klalb/LINESPacket.java +++ b/src/org/kne/cloud/network/klalb/ADDLINESPacket.java @@ -5,7 +5,7 @@ import java.io.DataOutput; import java.io.IOException; import java.net.Inet6Address; -public class LINESPacket extends KLALBPacket { +public class ADDLINESPacket extends KLALBPacket { private String lines; @@ -13,13 +13,13 @@ public class LINESPacket extends KLALBPacket { return lines; } - public LINESPacket(String lines) { - super(LINES); + public ADDLINESPacket(String lines) { + super(ADDLINES); this.lines=lines; } - public LINESPacket() { - super(LINES); + public ADDLINESPacket() { + super(ADDLINES); } private int UTFlength(String str) { int strlen = str.length(); @@ -38,7 +38,7 @@ public class LINESPacket extends KLALBPacket { } @Override public String toString() { - return "LINES\n"+lines; + return "ADDLINES\n"+lines; } @Override diff --git a/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java b/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java index c4490f6..0516448 100644 --- a/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java +++ b/src/org/kne/cloud/network/klalb/ByteArrayRecycle.java @@ -11,14 +11,14 @@ public class ByteArrayRecycle { this.capacity = capacity; this.length = length; } - public synchronized void recycle(byte[]b) { + public void recycle(byte[]b) { if(b.length!=length) throw new IllegalArgumentException("wrong length"); if(rec.size(),PortPacket{ + public static final ByteArrayRecycle arrayRecycle=new ByteArrayRecycle(5000,65535); @Override public long getLength() { - return super.getLength()+18+size; + return super.getLength()+19+size; } private int sport,dport; private long number; private int size; private byte[]data; + private int sendcount; + + volatile long resendtimer=System.nanoTime(); + public DATATPacket(int sport,int dport,long number,byte[]data,int size) { super(DATAT); this.sport=sport; @@ -24,6 +34,10 @@ public class DATATPacket extends KLALBPacket { this.size=size; } + public int getSendcount() { + return sendcount; + } + public DATATPacket() { super(DATAT); } @@ -48,28 +62,107 @@ public class DATATPacket extends KLALBPacket { public String toString() { return "DATAT "+sport+"->"+dport+" "+number+"["+size+"]"; } - + private final byte[] writeBuffer = new byte[19]; @Override protected void writeToStream(DataOutput dto) throws IOException { super.writeToStream(dto); - dto.writeInt(sport); + /*dto.writeInt(sport); dto.writeInt(dport); dto.writeLong(number); - dto.writeChar(size); + dto.writeByte(getSendRecord().size()); + dto.writeChar(size);*/ + + writeBuffer[0] = (byte)(sport >>> 24); + writeBuffer[1] = (byte)(sport >>> 16); + writeBuffer[2] = (byte)(sport >>> 8); + writeBuffer[3] = (byte)(sport >>> 0); + + writeBuffer[4] = (byte)(dport >>> 24); + writeBuffer[5] = (byte)(dport >>> 16); + writeBuffer[6] = (byte)(dport >>> 8); + writeBuffer[7] = (byte)(dport >>> 0); + + writeBuffer[8] = (byte)(number >>> 56); + writeBuffer[9] = (byte)(number >>> 48); + writeBuffer[10] = (byte)(number >>> 40); + writeBuffer[11] = (byte)(number >>> 32); + writeBuffer[12] = (byte)(number >>> 24); + writeBuffer[13] = (byte)(number >>> 16); + writeBuffer[14] = (byte)(number >>> 8); + writeBuffer[15] = (byte)(number >>> 0); + + writeBuffer[16]=(byte) getSendRecord().size(); + + writeBuffer[17]=(byte) (size>>>8); + writeBuffer[18]=(byte) (size>>>0); + dto.write(writeBuffer); + //CRC32 crc=new CRC32(); + // crc.update(data, 0, size); dto.write(data,0,size); + //dto.writeLong(crc.getValue()); } + private final byte[] readBuffer = new byte[19]; @Override protected void readFromStream(DataInput din) throws IOException { super.readFromStream(din); - sport=din.readInt(); + /*sport=din.readInt(); dport=din.readInt(); number=din.readLong(); - size=din.readChar(); + sendcount=din.readUnsignedByte(); + size=din.readChar();*/ + + din.readFully(readBuffer); + sport=((readBuffer[0] << 24) + (readBuffer[1] << 16) + (readBuffer[2] << 8) + (readBuffer[3] << 0)); + dport=((readBuffer[4] << 24) + (readBuffer[5] << 16) + (readBuffer[6] << 8) + (readBuffer[7] << 0)); + number=(((long)readBuffer[8] << 56) + + ((long)(readBuffer[9] & 255) << 48) + + ((long)(readBuffer[10] & 255) << 40) + + ((long)(readBuffer[11] & 255) << 32) + + ((long)(readBuffer[12] & 255) << 24) + + ((readBuffer[13] & 255) << 16) + + ((readBuffer[14] & 255) << 8) + + ((readBuffer[15] & 255) << 0)); + sendcount=readBuffer[16]&0xff; + size=(((readBuffer[17]&0xff) << 8) + ((readBuffer[18]&0xff) << 0)); + data=arrayRecycle.create(); din.readFully(data,0,size); + /*CRC32 crc32=new CRC32(); + crc32.update(data, 0, size); + if(din.readLong()!=crc32.getValue()) { + throw new StreamCorruptedException("CRC32 error!"); + }*/ } + @Override + public int hashCode() { + return Objects.hash(number); + } + + @Override + public boolean equals(Object obj) { + if (this == obj) + return true; + if (obj == null) + return false; + if (getClass() != obj.getClass()) + return false; + DATATPacket other = (DATATPacket) obj; + return number == other.number; + } + public int getSize() { return size; } + + @Override + public int compareTo(DATATPacket o) { + if(number>o.number) { + return 1; + }else if(number selflineTableSupplier=()->{return null;}; - - public Supplier getSelflineTableSupplier() { return selflineTableSupplier; @@ -53,149 +75,156 @@ public class KLALBController { this.selflineTableSupplier = selflineTableSupplier; } - public LineManager getLineManager() { - if(lineManager==null) - lineManager=createLineManager(); - return lineManager; - } - - protected LineManager createLineManager() { - return new LineManager(this); - } - private Inet6Address self; public Inet6Address getSelf() { return self; } - - private Map> routes = new ConcurrentHashMap<>(); - - protected Map> getRoutes() { - return routes; + private List lines=new ArrayList<>(); + //private ReadWriteLock lineslock=new ReentrantReadWriteLock(); + + public List getLines() { + return lines; } - private void setLine(Inet6Address vaddr, KLALBRemoteSocket krs) { - synchronized (routes) { - List al = routes.computeIfAbsent(vaddr, (vaddr2) -> { - return new ArrayList(); - }); - al.add(krs); - } + private PortBinder streamPortBinder=new PortBinder(this); + + protected PortBinder getStreamPortBinder() { + return streamPortBinder; } - private void removeLine(KLALBRemoteSocket krs) { - synchronized (routes) { - Iterator>> iter = routes.entrySet().iterator(); - while (iter.hasNext()) { - Map.Entry> entry = (Map.Entry>) iter - .next(); - entry.getValue().remove(krs); - if (entry.getValue().isEmpty()) { - iter.remove(); - } + + public void reconnectImmediately() { + synchronized (lines) { + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + klalbRemoteLine.reconnectImmediately(); } } } + + private PacketReceiver prc=new PacketReceiver(); + private class PacketReceiver implements KLALBPacketConsumer{ - public void addRemoteSocket(KLALBRemoteSocket krs) { - krs.setController(this); - CountDownLatch cdl = new CountDownLatch(1); - krs.setPacketReceiver((rec) -> { + @Override + public void accept(KLALBRemoteLine krs, KLALBPacket rec) { try { - if (rec instanceof SYNTPacket) { - SYNTPacket synt = (SYNTPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(synt.getDport()); - if (kvi != null) { - if (kvi.isListening()) { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), synt); - } - } else { - - sendPacketToAddress(krs.getRemoteVaddr(), new RSTPacket(synt.getDport(), synt.getSport()), - 65537); + if(rec instanceof PortPacket&&krs.getRemoteVaddr()!=null) { + PortPacket pt=(PortPacket) rec; + if(!streamPortBinder.distributePacketToConsumer(krs, pt)) { + if(!(pt instanceof RSTPacket)) + sendPacketToAddress(krs.getRemoteVaddr(), new RSTPacket(pt.getDport(), pt.getSport()), + 0,2); } - } else if (rec instanceof SACKTPacket) { - SACKTPacket sackt = (SACKTPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(sackt.getDport()); - if (kvi != null) { - if (!kvi.isListening()) { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), sackt); - } - } - } else if (rec instanceof RSTPacket) { - RSTPacket rst = (RSTPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(rst.getDport()); - if (kvi != null) { - if (kvi.isListening()) { - KLALBVirtualSocketImpl kvi2 = kvi.getAccepts() - .get(new InetSocketAddress(krs.getRemoteVaddr(), rst.getSport())); - if (kvi2 != null) { - kvi2.getPackReceiver().accept(krs.getRemoteVaddr(), rst); - } - } else { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), rst); - } - } - } else if (rec instanceof DATATPacket) { - DATATPacket datat = (DATATPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(datat.getDport()); - if (kvi != null) { - if (kvi.isListening()) { - KLALBVirtualSocketImpl kvi2 = kvi.getAccepts() - .get(new InetSocketAddress(krs.getRemoteVaddr(), datat.getSport())); - if (kvi2 != null) { - kvi2.getPackReceiver().accept(krs.getRemoteVaddr(), datat); - } - } else { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), datat); - } - } - } else if (rec instanceof ACKTPacket) { - ACKTPacket ackt = (ACKTPacket) rec; - KLALBVirtualSocketImpl kvi = bindmap.get(ackt.getDport()); - if (kvi != null) { - if (kvi.isListening()) { - KLALBVirtualSocketImpl kvi2 = kvi.getAccepts() - .get(new InetSocketAddress(krs.getRemoteVaddr(), ackt.getSport())); - if (kvi2 != null) { - kvi2.getPackReceiver().accept(krs.getRemoteVaddr(), ackt); - } - } else { - kvi.getPackReceiver().accept(krs.getRemoteVaddr(), ackt); - } - } - } else if (rec instanceof VADDRPacket) { - VADDRPacket var = (VADDRPacket) rec; - setLine(var.getVaddr(), krs); - cdl.countDown(); - }else if(rec instanceof LINESPacket) { - LINESPacket lpt=(LINESPacket) rec; + }else { + switch (rec.getType()) { + case KLALBPacket.ADDLINES: + ADDLINESPacket lpt=(ADDLINESPacket) rec; String s=lpt.getLines(); Scanner scn=new Scanner(s); while(scn.hasNext()) { String sn=scn.nextLine(); - getLineManager().addHostPort(new MultipurposeSocketAddress(sn)); + MultipurposeSocketAddress msa= new MultipurposeSocketAddress(sn); + addRemoteLines(msa); + } + break; + default: + break; } - } catch (NoRouteToHostException e) { + } + } catch (IOException e) { + // TODO 自动生成的 catch 块 e.printStackTrace(); } - }); - krs.setCloseListener((x) -> { - removeLine(krs); - cdl.countDown(); - }); - krs.sendPacket(new VADDRPacket(self), 65537); - String selflineTable=selflineTableSupplier.get(); - if(selflineTable!=null) - krs.sendPacket(new LINESPacket(selflineTable), 65537); - try { - cdl.await(); - } catch (InterruptedException e) { - e.printStackTrace(); + } + + + } + public void addRemoteLines(MultipurposeSocketAddress target) throws SocketTimeoutException, SocketException { + synchronized (lines) { + + try { + Enumerationeu= NetworkInterface.getNetworkInterfaces(); + while (eu.hasMoreElements()) { + NetworkInterface networkInterface = (NetworkInterface) eu.nextElement(); + if(networkInterface.isUp()) { + //System.out.println(networkInterface+" "+networkInterface.isUp()); + Enumerationei= networkInterface.getInetAddresses(); + while (ei.hasMoreElements()) { + InetAddress inetAddress = (InetAddress) ei.nextElement(); + MultipurposeSocketAddress bind=new MultipurposeSocketAddress(inetAddress.getHostAddress(),0); + if(!checkContainsTargetAndBind(target,bind)) { + //System.out.println(target+" "+bind); + addRemoteLine( new KLALBRemoteLine(target,bind)); + } + } + } + } + } catch (SocketException e) { + if(!checkContainsTarget(target)) + addRemoteLine( new KLALBRemoteLine(target)); + throw e; + } + } } + public Inet6Address getRemoteVaddrBySocketAddress(MultipurposeSocketAddress target) throws SocketTimeoutException { + KLALBRemoteLine kr=null; + synchronized(lines) { + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + if(target.equals(klalbRemoteLine.getSocketAddress())) { + kr=klalbRemoteLine; + } + } + } + if(kr==null) { + kr=new KLALBRemoteLine(target); + addRemoteLine(kr); + kr.waitForRemoteVaddrAvaliable(20000); + }else { + kr.reconnectImmediately(); + kr.waitForRemoteVaddrAvaliable(20000); + } + return kr.getRemoteVaddr(); + } + private boolean checkContainsTargetAndBind(MultipurposeSocketAddress target,MultipurposeSocketAddress bind) { + boolean b=false; + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + if(bind.equals(klalbRemoteLine.getBindAddress())&&target.equals(klalbRemoteLine.getSocketAddress())) { + b=true; + break; + } + } + return b; + } + + private boolean checkContainsTarget(MultipurposeSocketAddress target) { + boolean b=false; + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + if(target.equals(klalbRemoteLine.getSocketAddress())) { + b=true; + break; + } + } + return b; + } + + public void addRemoteLine(KLALBRemoteLine krs) throws SocketTimeoutException { + krs.setPacketReceiver(prc); + krs.setLocalVaddrSupplier(()->{return self;}); + krs.startIO(); + String selflineTable=selflineTableSupplier.get(); + if(selflineTable!=null) + krs.sendPacket(new ADDLINESPacket(selflineTable), 0); + synchronized (lines) { + lines.add(krs); + } + } + public KLALBController(Inet6Address self) { this.self = self; @@ -205,135 +234,149 @@ public class KLALBController { this.self=KLALBUtils.uuidToIP(UUID.randomUUID()); } - private Map bindmap = new ConcurrentHashMap<>(); - - protected Map getBindmap() { - return bindmap; - } - protected KLALBVirtualSocketImpl createVirtualImpl() { return new KLALBVirtualSocketImpl(this); } - protected int bind(KLALBVirtualSocketImpl klalbVirtualSocketImpl, int port) throws BindException { - synchronized (bindmap) { - if (port == 0) { - port = allocPort(); - } - if (bindmap.putIfAbsent(port, klalbVirtualSocketImpl) != null) { - throw new BindException("port " + port + " is already bind!"); - } - return port; - } - } - - protected void unbind(KLALBVirtualSocketImpl klalbVirtualSocketImpl) { - synchronized (bindmap) { - Set> s = bindmap.entrySet(); - Iterator> it = s.iterator(); - while (it.hasNext()) { - Entry object = it.next(); - if (klalbVirtualSocketImpl.equals(object.getValue())) { - it.remove(); - return; - } - } - } - } - - protected int allocPort() throws BindException { - int i = 1; - while (bindmap.containsKey(i)) { - if (i == 65535) { - throw new BindException("can't alloc port"); - } - i++; - } - return i; - } + protected void sendPacketToAddress(Inet6Address addr, KLALBPacket syntPacket, int priority) - throws NoRouteToHostException { + throws IOException { sendPacketToAddress(addr, syntPacket, priority, 1); } - - protected void sendPacketToAddress(Inet6Address addr, KLALBPacket packet, int priority, int count) - throws NoRouteToHostException { - synchronized (routes) { - List l = routes.get(addr); - if (l == null || l.isEmpty()) { - throw new NoRouteToHostException("address unreachable: " + addr); + + private void updateLines2(Inet6Address addr) throws SocketTimeoutException { + List l=new ArrayList(); + synchronized (lines) { + for (int i = 0; i < lines.size(); i++) { + KLALBRemoteLine klalbRemoteLine = lines.get(i); + if(klalbRemoteLine.isClosed()) { + lines.remove(i); + i--; + continue; + } + if(addr.equals(klalbRemoteLine.getRemoteVaddr())&&klalbRemoteLine.getMonitor().getState()==Monitor.ONLINE) { + l.add(klalbRemoteLine); + } + } + } + if(l.isEmpty()) { + lines2.remove(addr); + }else { + lines2.put(addr, l); + } } - List l2 = (List) ((ArrayList) l).clone(); + private Map> lines2 = new ConcurrentHashMap<>(); + private volatile long itm=System.nanoTime(); + protected void sendPacketToAddress(Inet6Address addr, KLALBPacket packet, int priority, int count) + throws IOException { + /*if(packet instanceof RSTPacket) { + new Exception("-RST-").printStackTrace(); + }*/ + //TimeDebugger tdb=new TimeDebugger(); + //tdb.putTime("start"); + + List lines2x; + //loop:while(true) { + long cur=System.nanoTime(); + if(cur-itm>10000000L) { + itm=cur; + lines2.clear(); + } + lines2x=lines2.get(addr); + if (lines2x == null || lines2x.isEmpty()) { + updateLines2(addr); + lines2x=lines2.get(addr); + } + if (lines2x == null || lines2x.isEmpty()) { + throw new NoRouteToHostException("address unreachable: " + addr); + } + //tdb.putTime("selectLines"); + /* for (Iterator iterator = lines2x.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + if(klalbRemoteLine.statLengthBefore(priority)<=65536*10) { + break loop; + } + } + try { + //System.out.println("slp"); + Thread.sleep(1); + } catch (InterruptedException e) { + e.printStackTrace(); + } + }*/ + List l2 = (List) ((ArrayList) lines2x).clone(); l2.removeAll(packet.getSendRecord()); if(l2.isEmpty()) { - l2 = (List) ((ArrayList) l).clone(); + l2 = (List) ((ArrayList) lines2x).clone(); } + + //tdb.putTime("findAvaliable"); + int count0 = Math.min(count, l2.size()); - LineDecitionComparator ldc = new LineDecitionComparator(l2,packet, priority); - Collections.sort(l2,ldc); + l2.forEach((r)->{ + r.runPredict(packet,priority); + }); + //Collections.shuffle(l2); + Collections.sort(l2); + //tdb.putTime("makeDecision"); //System.out.println(l2); - for (Iterator iterator = l2.iterator(); iterator.hasNext();) { - KLALBRemoteSocket krst = (KLALBRemoteSocket) iterator.next(); + for (int i = 0; i < l2.size(); i++) { + KLALBRemoteLine krst =l2.get(i); krst.sendPacket(packet, priority); packet.getSendRecord().add(krst); count0--; - Thread.yield(); if (count0 <= 0) break; } - } - + //tdb.putTime("sendPacket"); + //tdb.print(); + /*if(packet instanceof RSTPacket) + new Exception().printStackTrace();*/ } protected void removeFromSend(Inet6Address addr,KLALBPacket klalbPacket) { - synchronized (routes) { - List l = routes.get(addr); - if (l != null && !l.isEmpty()) { - l.forEach((x)->{ - x.remoeFromSendQueue(klalbPacket); - }); + synchronized (lines) { + for (Iterator iterator = lines.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteLine = (KLALBRemoteLine) iterator.next(); + klalbRemoteLine.remoeFromSendQueue(klalbPacket); } } } - protected boolean checkIsBind(KLALBVirtualSocketImpl klalbVirtualSocketImpl) { - return bindmap.containsValue(klalbVirtualSocketImpl); - } private Timer t=new Timer("数据包发送计时器", true); public Timer getTimer() { return t; } - - public void registerToProxyTypeAs(String proxyname) { - MultipurposeSocketAddress.getSocketFactoryRegister().put(proxyname, new KLALBVirtualSocketFactory(this)); - MultipurposeSocketAddress.getServerSocketFactoryRegister().put(proxyname, new KLALBVirtualServerSocketFactory(this)); - ProxyProfileEntry.getRegister().put(proxyname, new KLALBNetworkService()); - } - private class KLALBNetworkService implements NetworkService{ + + public void registerToProxyTypeAs(String proxyname) { + MultipurposeSocketAddress.getSocketTypeRegister().put(proxyname+"_Stream",socketType); + //ProxyProfileEntry.getRegister().put(proxyname, this); + } - @Override - public void listen(MultipurposeSocketAddress msa) { - // TODO 自动生成的方法存根 - - } - @Override - public void unlisten(MultipurposeSocketAddress msc) { - // TODO 自动生成的方法存根 - - } - - @Override - public void connect(MultipurposeSocketAddress msa) { - getLineManager().addHostPort(msa); - } - - @Override - public void unconnect(MultipurposeSocketAddress msc) { - getLineManager().removeHostPort(msc); - } + /*@Override + public void listen(MultipurposeSocketAddress msa) { + // TODO 自动生成的方法存根 } + @Override + public void unlisten(MultipurposeSocketAddress msc) { + // TODO 自动生成的方法存根 + + } + + @Override + public void connect(MultipurposeSocketAddress msa) { + // TODO 自动生成的方法存根 + + } + + @Override + public void unconnect(MultipurposeSocketAddress msc) { + // TODO 自动生成的方法存根 + + }*/ + } diff --git a/src/org/kne/cloud/network/klalb/KLALBInputStream.java b/src/org/kne/cloud/network/klalb/KLALBInputStream.java index af6489b..915a38e 100644 --- a/src/org/kne/cloud/network/klalb/KLALBInputStream.java +++ b/src/org/kne/cloud/network/klalb/KLALBInputStream.java @@ -1,5 +1,7 @@ package org.kne.cloud.network.klalb; import static org.kne.cloud.network.klalb.KLALBPacket.*; + +import java.io.DataInput; import java.io.DataInputStream; import java.io.EOFException; import java.io.IOException; @@ -21,49 +23,68 @@ public class KLALBInputStream extends DataInputStream { if(bv!=2) throw new StreamCorruptedException("remote version is V"+bv+"."+sv+",not V2.0"); } - public synchronized KLALBPacket readPacket() throws IOException { - int type=read(); + public KLALBPacket readPacket() throws IOException { + return readKLALBPacketFromStream(this); + } + public static KLALBPacket readKLALBPacketFromStream(DataInputStream in) throws IOException { +int type=in.read(); if(type==-1) { - throw new EOFException(); + return null; } KLALBPacket klp; switch(type) { case PING:klp=new PINGPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case PONG: klp=new PONGPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case SYNT: klp=new SYNTPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case SACKT: klp=new SACKTPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case RST: klp=new RSTPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case DATAT: klp=new DATATPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case ACKT: klp=new ACKTPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; case VADDR: klp=new VADDRPacket(); - klp.readFromStream(this); + klp.readFromStream(in); return klp; - case LINES: - klp=new LINESPacket(); - klp.readFromStream(this); + case ADDLINES: + klp=new ADDLINESPacket(); + klp.readFromStream(in); return klp; + case NACKT: + klp=new NACKTPacket(); + klp.readFromStream(in); + return klp; + case TEST: + klp=new TESTPacket(); + klp.readFromStream(in); + return klp; + case VADDRACK: + klp=new VADDRACKPacket(); + klp.readFromStream(in); + return klp; + case VADDRREQ: + klp=new VADDRREQPacket(); + klp.readFromStream(in); + return klp; } throw new StreamCorruptedException("unknown package type:"+type); } diff --git a/src/org/kne/cloud/network/klalb/KLALBMain.java b/src/org/kne/cloud/network/klalb/KLALBMain.java index b9440ab..7624060 100644 --- a/src/org/kne/cloud/network/klalb/KLALBMain.java +++ b/src/org/kne/cloud/network/klalb/KLALBMain.java @@ -1,72 +1,58 @@ package org.kne.cloud.network.klalb; +import java.io.BufferedReader; import java.io.File; +import java.io.FileNotFoundException; +import java.io.FileReader; import java.io.IOException; import java.util.Iterator; -import java.util.List; import java.util.Scanner; -import org.kne.cloud.network.klalb.LineManager.LineEntry; -import org.kne.cloud.network.mport.MultipurposeSocketAddress; -import org.kne.cloud.network.mport.ProxyProfileAnalyser; -import org.kne.cloud.network.mport.ProxyProfileExecutor; -import org.kne.cloud.network.mport.TCPListener; +import org.kne.cloud.network.KLALBProxyConfigJsonExecuter; +import org.kne.cloud.network.SocketToServiceProxy; +import org.kne.cloud.network.SocketToSocketProxy; public class KLALBMain { - public static TCPListener tcpl; - public static ServerPropties sp; - public static KLALBController kc; - - public static ProxyProfileExecutor pfa; - public static File f=new File("lines.cfg"); + public static void main(String[] args) throws IOException { - - System.out.println("KLALB负载均衡V2.0"); - sp=new ServerPropties(); - - System.out.println("虚拟地址:"+sp.getVirtualIP().getHostAddress()); - kc=new KLALBController(sp.getVirtualIP()); - kc.registerToProxyTypeAs("KLALB"); - - - - - pfa=new ProxyProfileExecutor(); - pfa.load(f); - + System.out.println(CONST.klalb+" V"+CONST.klalbver); Scanner scn=new Scanner(System.in); + + File configJson=new File("klalbconfig.json"); + + KLALBProxyConfigJsonExecuter kpcje=new KLALBProxyConfigJsonExecuter(); + kpcje.loadConfigJson(configJson); while(true) { String s=scn.next(); String[]sc=s.split(" "); switch(sc[0]) { case "help": + System.out.print("help:查看命令使用说明"); System.out.print("state:查看线路状态"); - System.out.println("reload:重新加载线路配置"); + //System.out.println("reload:重新加载线路配置文件"); + System.out.println("reconnect:所有离线线路跳过重连等待时间立即尝试重连"); + System.out.println("stop:退出程序"); + break; - case "reloadlines": - pfa.load(f); - System.out.println("重新加载线路配置成功"); - break; case "state": - LineManager le=kc.getLineManager(); - for (Iterator> iterator = le.entrySet().iterator(); iterator.hasNext();) { - java.util.Map.Entry hostPort = iterator.next(); - System.out.println(hostPort.getValue() .toString()); + System.out.println("状态\t上传流量\t下载流量\t上传速度\t下载速度\t延迟\t抖动\t下一次重试"); + for (Iterator iterator = kpcje.getKlalbController().getLines().iterator(); iterator.hasNext();) { + KLALBRemoteLine hostPort = iterator.next(); + System.out.println(hostPort .toString2()); + //System.out.println(); } break; + case "stop": + System.exit(0); + break; + case "reconnect": + kpcje.getKlalbController().reconnectImmediately(); + break; default: System.out.println("未知命令,请输入help以查询指令说明"); } } + } - private static void openPort(int port) throws IOException { - if(tcpl!=null) - tcpl.close(); - tcpl=new TCPListener(new MultipurposeSocketAddress("0.0.0.0", port)); - tcpl.setCon((soc)->{ - KLALBRemoteSocket krs=new KLALBRemoteSocket(soc); - kc.addRemoteSocket(krs); - }); - tcpl.open(); - } + } diff --git a/src/org/kne/cloud/network/klalb/KLALBOutputStream.java b/src/org/kne/cloud/network/klalb/KLALBOutputStream.java index 3de1e53..457b344 100644 --- a/src/org/kne/cloud/network/klalb/KLALBOutputStream.java +++ b/src/org/kne/cloud/network/klalb/KLALBOutputStream.java @@ -12,11 +12,15 @@ public class KLALBOutputStream extends DataOutputStream { byte[]b=new byte[] {'K','L','A','L','B'}; write(b); writeInt(2); - writeInt(0); + writeInt(1); flush(); } - public synchronized void writePacket(KLALBPacket klb) throws IOException { - write(klb.getType()); - klb.writeToStream(this); + public void writePacket(KLALBPacket klb) throws IOException { + writeKLALBPacketToStream(this, klb); + } + public static void writeKLALBPacketToStream(DataOutputStream out,KLALBPacket klb) throws IOException { + out.write(klb.getType()); + klb.writeToStream(out); + out.flush(); } } diff --git a/src/org/kne/cloud/network/klalb/KLALBPacket.java b/src/org/kne/cloud/network/klalb/KLALBPacket.java index 8bc5ac0..5f65c2d 100644 --- a/src/org/kne/cloud/network/klalb/KLALBPacket.java +++ b/src/org/kne/cloud/network/klalb/KLALBPacket.java @@ -8,7 +8,7 @@ import java.io.ObjectInput; import java.io.ObjectOutput; import java.util.Vector; -public abstract class KLALBPacket{ +public abstract class KLALBPacket implements Sumable{ public static final int PING=0; public static final int PONG=1; public static final int SYNT=2; @@ -18,7 +18,11 @@ public abstract class KLALBPacket{ public static final int ACKT=6; public static final int DATAU=7; public static final int VADDR=8; - public static final int LINES=9; + public static final int ADDLINES=9; + public static final int NACKT=10; + public static final int TEST=11; + public static final int VADDRACK=12; + public static final int VADDRREQ=13; private int type; public KLALBPacket(int type) { @@ -32,18 +36,32 @@ public abstract class KLALBPacket{ public int getType() { return type; } + + private long sndtime,rcvtime; protected void writeToStream(DataOutput dto) throws IOException { - + sndtime=System.nanoTime(); } protected void readFromStream(DataInput din) throws IOException { - + rcvtime=System.nanoTime(); + } + public long getSndtime() { + return sndtime; + } + public long getRcvtime() { + return rcvtime; } public long getLength() { return 1; } - private Vector sendRecord=new Vector<>(); + @Override + public long getValue() { + return getLength(); + } - public Vector getSendRecord() { + + private Vector sendRecord=new Vector<>(); + + public Vector getSendRecord() { return sendRecord; } diff --git a/src/org/kne/cloud/network/klalb/KLALBUtils.java b/src/org/kne/cloud/network/klalb/KLALBUtils.java index 3e4a5bb..f0e314f 100644 --- a/src/org/kne/cloud/network/klalb/KLALBUtils.java +++ b/src/org/kne/cloud/network/klalb/KLALBUtils.java @@ -5,6 +5,7 @@ import java.io.DataOutputStream; import java.io.IOException; import java.net.Inet6Address; import java.net.UnknownHostException; +import java.util.List; import java.util.UUID; public class KLALBUtils { @@ -54,4 +55,24 @@ public class KLALBUtils { } return fmt; } + public static int binarySearchDATATPacketNumber(Listlst,long number) { + int low =0; + int high=lst.size()-1; + int middle=0; + if(low>high||numberlst.get(high).getNumber()) { + return -1; + } + while(low<=high) { + middle=(low+high)/2; + long xn=lst.get(middle).getNumber(); + if(xn>number) { + high=middle-1; + }else if(xn backlogQueue; + //private CountDownLatch reseted=new CountDownLatch(1); + private BlockingQueue backlogQueue; + private Thread sendDequeLock; private Queue sendDeque = new ConcurrentLinkedQueue<>(); - private List sendlist=new Vector<>(); + private Map sendmap=new HashMap<>(); + private volatile Thread sendthread; + //private List sendlist=new ArrayList<>(); private TimerTask sendCheckTask=new SendCheckTask(); private class SendCheckTask extends TimerTask{ - + public void run() { - synchronized (sendlist) { - for (int i = 0; i < sendlist.size(); i++) { + synchronized (sendmap) { + Collection cdp=sendmap.values(); + for (Iterator iterator = cdp.iterator(); iterator.hasNext();) { + DATATPacket dtp = (DATATPacket) iterator.next(); try { - sendlist.get(i).check(i); + long x=System.nanoTime(); + long dt=x-dtp.resendtimer; + long limit= dtp.getSendRecord().size()*RTTAvg*6+10000000; + if(dt>limit) { + if(dtp.getSendRecord().size()>=20) { + throw new IOException("send error!"); + } + controller.sendPacketToAddress((Inet6Address) remoteaddr,dtp, 5-1); + System.out.println("第"+(dtp.getSendRecord().size()-1)+"次重传:"+dtp); + dtp.resendtimer=x; + + } + } catch (IOException e) { + e.printStackTrace(); + try { + close0(true); + } catch (IOException e1) { + e1.printStackTrace(); + } + break; + } + } + + } + /*for (int i = 0; i < sendlist.size(); i++) { + try { + DATATPacket dtp= sendlist.get(i); + long x=System.nanoTime(); + long dt=x-dtp.resendtimer; + long limit= dtp.getSendRecord().size()*(1000000000L*i+RTTAvg); + if(dt>limit) { + if(dtp.getSendRecord().size()>=30) { + throw new IOException("send error!"); + } + controller.sendPacketToAddress((Inet6Address) address,dtp, 5-1); + System.out.println("第"+(dtp.getSendRecord().size()-1)+"次重传:"+dtp); + dtp.resendtimer=x; + + } } catch (IOException e) { e.printStackTrace(); try { @@ -101,111 +173,56 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } break; } - } - } + }*/ + + //System.out.println(isClosed()); } } - private volatile boolean succeed, refused; + + private TimerTask flowControlTask=new FlowControlTask(); + private class FlowControlTask extends TimerTask{ + @Override + public void run() { + if(sendDeque.isEmpty()) { + try {if(getLocalPort()!=0&&getPort()!=0) + if(remoteaddr instanceof Inet6Address&&(!remoteaddr.isAnyLocalAddress())) + controller.sendPacketToAddress((Inet6Address) remoteaddr, new ACKTPacket(getLocalPort(),getPort(), -1, + true,0), 0); + } catch (IOException e) { + try { + close0(true); + } catch (IOException e1) { + e1.printStackTrace(); + } + e.printStackTrace(); + } + + } + } + + } + + protected boolean isListening() { return backlogQueue != null; } - private List inputchache = new ArrayList<>(); + //private int inputcross = 0; + private List inputchache = new ArrayList<>(); private long inputcount = 0; private volatile boolean avaliable = true; - private BiConsumer packReceiver = new BiConsumer() { + + + private long RTTAvg=1000000000L; + private volatile boolean ignoreBindCheck; + private boolean connected; - @Override - public void accept(Inet6Address from, KLALBPacket u) { - try { - if (u instanceof SYNTPacket) { - if (isListening()) { - if (backlogQueue.offer(new InetSocketAddress(from, ((SYNTPacket) u).getSport()))) { - controller.sendPacketToAddress(from, - new SACKTPacket(localport, ((SYNTPacket) u).getSport()), 65537,2); - - } else { - controller.sendPacketToAddress(from, new RSTPacket(localport, ((SYNTPacket) u).getSport()), - 65537,2); - } - } else { - controller.sendPacketToAddress(from, new RSTPacket(localport, ((SYNTPacket) u).getSport()), - 65537,2); - } - } else if (u instanceof SACKTPacket) { - if (connecting) { - connecting = false; - succeed = true; - cdl.countDown(); - } - } else if (u instanceof RSTPacket) { - if (connecting) { - connecting = false; - refused = true; - cdl.countDown(); - close(); - } else { - close(); - } - } else if (u instanceof DATATPacket) { - DATATPacket dtp = (DATATPacket) u; - synchronized (inputchache) { - if (dtp.getNumber() >= inputcount) { - inputchache.add(dtp); - while (true) { - DATATPacket kkb = null; - for (int i = 0; i < inputchache.size(); i++) { - DATATPacket klalbBlock = inputchache.get(i); - if (klalbBlock.getNumber() == inputcount) { - inputchache.remove(i); - i--; - kkb = klalbBlock; - break; - } - } - 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,1); - } else if (u instanceof ACKTPacket) { - ACKTPacket ackt = (ACKTPacket) u; - avaliable = ackt.isAvaliable(); - 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(); - } - } - - - - }; + + + private int countInputBytes() { AtomicInteger i = new AtomicInteger(0); sendDeque.forEach((c) -> { @@ -215,23 +232,17 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } private int countOutputBytes() { AtomicInteger i = new AtomicInteger(0); - sendlist.forEach((c) -> { - i.addAndGet(c.getPacket().getSize()); + sendmap.values().forEach((c) -> { + i.addAndGet(c.getSize()); }); return i.get(); } public KLALBController getController() { return controller; } - - public BiConsumer getPackReceiver() { - return packReceiver; - } - public KLALBVirtualSocketImpl(KLALBController kc) { super(); this.controller = kc; - kc.getTimer().schedule(sendCheckTask, 50, 50); } @Override @@ -259,7 +270,7 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { case SocketOptions.SO_SNDBUF: return outputchachesize; case SocketOptions.SO_BINDADDR: - return bindaddr; + return localaddr; default: return null; } @@ -283,43 +294,58 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { connect(new InetSocketAddress(address, port), 10000); } - private CountDownLatch cdl = new CountDownLatch(1); - @Override protected void connect(SocketAddress address, int timeout) throws IOException { - if (!controller.checkIsBind(this)) { + if(!ignoreBindCheck) + if (!controller.getStreamPortBinder().checkIsBind(this)) { bind(Inet6Address.getByName("::0"), 0); } - connecting = true; port = ((InetSocketAddress) address).getPort(); - this.address = ((InetSocketAddress) address).getAddress(); - controller.sendPacketToAddress((Inet6Address) this.address, new SYNTPacket(localport, port), 65537); + this.address=this.remoteaddr = (Inet6Address) ((InetSocketAddress) address).getAddress(); + controller.getStreamPortBinder().connect(this); + + + if(connected) + throw new SocketException("already connected"); + connected=true; + + controller.getTimer().schedule(sendCheckTask, 50, 50); + //controller.sendPacketToAddress((Inet6Address) this.remoteaddr, new SYNTPacket(localport, port), 0,1); try { - if (timeout == 0) { - cdl.await(); - } else { - cdl.await(timeout, TimeUnit.MILLISECONDS); + getKVSIOutputStream().write(compress); + getKVSIOutputStream().forceFlush(); + getKVSIOutputStream().waitForAllAcknowledged(timeout); + if(compress==1) { + vout=new DeflaterOutputStream(vout, new Deflater(Deflater.BEST_COMPRESSION, true), LIMIT, true); + }else { + vout = getKVSIOutputStream(); } - } catch (InterruptedException e) { - e.printStackTrace(); + int comp=getKVSIInputStream().read(); + if(comp==1) { + vin=new InflaterInputStream(getKVSIInputStream(),new Inflater(true),LIMIT); + }else { + vin = getKVSIInputStream(); + } - connecting = false; - if (succeed) { - } else if (refused) { - throw new ConnectException("connect refused"); - } else { + }catch(SocketTimeoutException e) { throw new SocketTimeoutException("connect time out"); + }catch(SocketException e) {//e.printStackTrace(); + throw new ConnectException("connect refused"); } + + controller.getTimer().schedule(flowControlTask, 5000, 5000); } - private volatile boolean connecting = false; + //private volatile boolean connecting = false; - public boolean isConnecting() { + /*public boolean isConnecting() { return connecting; - } - + }*/ @Override protected void bind(InetAddress host, int port) throws IOException { + bind(host,port,false); + } + protected void bind(InetAddress host, int port,boolean ignoreBindCheck) throws IOException { if (host.equals(Inet4Address.getByName("0.0.0.0"))) { host = Inet6Address.getByName("::0"); } @@ -329,46 +355,58 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { if ((!host.isAnyLocalAddress()) && (!host.equals(controller.getSelf()))) { throw new BindException("must bind to self"); } - bindaddr = (Inet6Address) host; - this.localport = controller.bind(this, port); + localaddr = (Inet6Address) host; + localport=port; + this.ignoreBindCheck=ignoreBindCheck; + if(!ignoreBindCheck) + controller.getStreamPortBinder().bind(this); + } + + @Override + public int getPort() { + return super.getPort(); } @Override protected void listen(int backlog) throws IOException { backlogQueue = new ArrayBlockingQueue<>(backlog); - address=bindaddr; - } - - private Map accepts = new ConcurrentHashMap<>(); - - public Map getAccepts() { - return accepts; + controller.getStreamPortBinder().listen(this); + address=localaddr; } @Override protected void accept(SocketImpl s) throws IOException { - try { KLALBVirtualSocketImpl kvsi = (KLALBVirtualSocketImpl) s; - InetSocketAddress isa = backlogQueue.take(); + + Object[] p=null; + while(true) { + if (isClosed()) + throw new SocketException("Socket is closed"); + p=backlogQueue.peek(); + if(p!=null) { + break; + } + try { + Thread.sleep(1); + } catch (InterruptedException e) { + e.printStackTrace(); + } + } + InetSocketAddress isa = (InetSocketAddress) p[0]; kvsi.inputchachesize=inputchachesize; kvsi.outputchachesize=outputchachesize; - kvsi.port = isa.getPort(); - kvsi.address = isa.getAddress(); + + kvsi.address = isa.getAddress(); + kvsi.remoteaddr=(Inet6Address) isa.getAddress(); + + kvsi.bind(localaddr, localport,true); + kvsi.accept((KLALBRemoteLine)p[1],(KLALBPacket) p[2]); + kvsi.connect(isa, 5000); + backlogQueue.poll(); + /*kvsi.port = isa.getPort(); kvsi.localport = localport; - kvsi.bindaddr = bindaddr; - accepts.put(isa, kvsi); - kvsi.setCloseListener((x) -> { - accepts.remove(isa); - }); - } catch (InterruptedException e) { - e.printStackTrace(); - } - } - - private Consumer acceptedSocketCloseListener; - - private void setCloseListener(Consumer lsr) { - this.acceptedSocketCloseListener = lsr; + kvsi.localaddr = localaddr;*/ + } private InputStream vin; @@ -382,27 +420,7 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { @Override public int read() throws IOException { 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(); - } - } + nextPacket(); } if (dtp.getSize() == 0) { return -1; @@ -416,6 +434,29 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } } + private DATATPacket nextPacket() throws IOException { + count = 0; + while (true) { + if (isClosed()) + throw new SocketException("Socket is closed"); + DATATPacket dtp2 = sendDeque.poll(); + if (dtp2 != null) { + dtp = dtp2; + checkFlowControl(dtp2); + break; + } + sendDequeLock=Thread.currentThread(); + LockSupport.parkNanos(1000000); + } + return dtp; + } + private void checkFlowControl(DATATPacket dtp2) throws IOException { + if(countInputBytes() >= inputchachesize-LIMIT*4) { + controller.sendPacketToAddress(remoteaddr, new ACKTPacket(dtp2.getDport(), dtp2.getSport(), dtp2.getNumber(), + true,dtp2.getSendcount()), 0); + } + } + @Override public int read(byte[] b, int off, int len) throws IOException { if (b == null) { @@ -430,29 +471,59 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { 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(); - } + nextPacket(); + } + if (dtp.getSize() == 0) { + return -1; + } else { + b[off]= dtp.getData()[count++] ; + if(count==dtp.getSize()) { + DATATPacket.arrayRecycle.recycle(dtp.getData()); + dtp=null; } } + int i = 1; + try { + while (i < len) { + + if (dtp == null ) { + nextPacket(); + } + if (dtp.getSize() == 0) { + break; + } + int remainD=len-i; + int min=Math.min(dtp.getSize()-count, remainD); + System.arraycopy(dtp.getData(), count, b, off + i, min); + count+=min; + i+=min; + //b[off + i]= dtp.getData()[count++] ; + if(count>=dtp.getSize()) { + DATATPacket.arrayRecycle.recycle(dtp.getData()); + dtp=null; + } + + } + } catch (IOException ee) { + } + return i; + } + /*@Override + public int read(byte[] b, int off, int len) throws IOException { + if (b == null) { + throw new NullPointerException(); + } else if (off < 0 || len < 0 || len > b.length - off) { + throw new IndexOutOfBoundsException(); + } else if (len == 0) { + return 0; + } + + len = Math.min(len, available()); + + + if (dtp == null ) { + nextPacket(); + } if (dtp.getSize() == 0) { return -1; } else { @@ -467,28 +538,7 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { for (; i < len; i++) { 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(); - } - } + nextPacket(); } if (dtp.getSize() == 0) { break; @@ -503,13 +553,18 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } catch (IOException ee) { } return i; - } - + }*/ @Override public void close() throws IOException { + try { + close0(); + }finally { + getKVSIOutputStream().close0(); + } + } + private void close0() throws IOException{ } - @Override public int available() throws IOException { AtomicInteger i = new AtomicInteger(0); @@ -524,138 +579,256 @@ public class KLALBVirtualSocketImpl extends VirtualSocketImpl { } private volatile long outputcount = 0; + private int compress=0; - private static final int LIMIT=65535; + public int getCompress() { + return compress; + } + + public void setCompress(int compress) { + this.compress = compress; + } + + private static final int LIMIT=8669; private class KVSIOutputStream extends OutputStream { + private byte[] cache=DATATPacket.arrayRecycle.create(); private int count=0; - private Object lock=new Object(); + private Lock olock=new ReentrantLock(); @Override public void write(int b) throws IOException { if (isClosed()) throw new SocketException("Socket is closed"); - synchronized (lock) { - - cache[count++]=(byte) b; + olock.lock(); + try { + cache[count++]=(byte) b;//1429?5773?8669 if (count >= LIMIT) {//1429?5773?8669 flush0(); }else { flush(); } + }finally { + olock.unlock(); } } + + @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) { + + olock.lock(); + try { + while(len>0) { + int remainD=LIMIT-count; + int min=Math.min(len, remainD); + System.arraycopy(b, off, cache, count, min); + count+=min; + off+=min; + len-=min; + if (count >= LIMIT) {//1429?5773?8669 + flush0(); + } + } + /*int ol=off+len; for (int i = off; i < ol; i++) { cache[count++]=b[i]; if (count >= LIMIT) {//1429?5773?8669 flush0(); } - } + }*/ flush(); + }finally { + olock.unlock(); + } + } + public void waitForAllAcknowledged(int timeout) throws IOException { + long start=System.nanoTime(); + while(true){ + if (isClosed()) + throw new SocketException("Socket is closed"); + //System.out.println(sendmap.size()); + if(sendmap.isEmpty()) + break; + if(timeout!=0&&(System.nanoTime()-start>timeout*1000000)) + throw new SocketTimeoutException("wait for acknowledged timout"); + sendthread=Thread.currentThread(); + LockSupport.parkNanos(1000000L); } } - 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(); + if(nodelay) { + olock.lock(); + try { + flush0(); + }finally { + olock.unlock(); + } + }else { + AtomicReference ioe=new AtomicReference<>(); + if(tt==null) { + tt=new TimerTask() { + + @Override + public void run() { + if(isClosed()) + cancel(); + try { + olock.lock(); + try { + flush0(); + }finally { + olock.unlock(); + } + } catch (IOException e) { + ioe.set(e); + } } - } catch (IOException e) { - e.printStackTrace(); + }; + new Timer("粘包计时线程").scheduleAtFixedRate(tt, delaytime, delaytime); + } + IOException ioex=ioe.get(); + if(ioex!=null) { + ioex.fillInStackTrace(); + throw ioex; } } - }; - new Timer("粘包计时线程").scheduleAtFixedRate(tt, 0, delaytime); - } + } + public void forceFlush() throws IOException{ + olock.lock(); + try { + flush0(); + }finally { + olock.unlock(); } } - private void flush0() throws IOException { if (count > 0) { + //TimeDebugger td=new TimeDebugger(); + //td.putTime("start"); + while (!avaliable) { try { - while (!avaliable) { Thread.sleep(1); - } } catch (InterruptedException e) { e.printStackTrace(); } - try { - while(countOutputBytes()>outputchachesize) { - synchronized (sendlist) { - sendlist.wait(10); } + //td.putTime("waitForAvaliable"); + while(true){ + if (isClosed()) + throw new SocketException("Socket is closed"); + //System.out.println(sendmap.size()); + boolean b=sendmap.size()<=outputchachesize/LIMIT; + if(b) + break; + sendthread=Thread.currentThread(); + LockSupport.parkNanos(1000000L); } - } catch (InterruptedException e) { - // TODO 自动生成的 catch 块 - e.printStackTrace(); - } + //td.putTime("waitForCache"); byte[] ba=cache; cache=DATATPacket.arrayRecycle.create(); - - SendTask x=new SendTask(controller, (Inet6Address) address, new DATATPacket(localport, port, outputcount++, ba,count), 5,20); + //td.putTime("flushBuffer"); + DATATPacket pack=new DATATPacket(localport, port, outputcount++, ba,count); + //td.putTime("createPacket"); + controller.sendPacketToAddress(remoteaddr,pack,1{ + try { + DATATPacket dp; + while((dp=kis.nextPacket()).getSize()!=0) { + los.write(dp.getData(), 0, dp.getSize()); + los.flush(); + } + }catch(IOException e) { + e.printStackTrace(); + }finally { + try { + los.close(); + } catch (IOException e) { + e.printStackTrace(); + } + try { + kis.close(); + } catch (IOException e) { + e.printStackTrace(); + } + } + }); + Thread t2=ThreadTool.makeVThreadIfSupport("本地接收线程", ()->{ + try { + int len=-1; + while(true) { + if((len=lis.read(kos.cache,0,LIMIT))==-1) { + break; + } + kos.olock.lock(); + try { + kos.count=len; + kos.flush0(); + }finally { + kos.olock.unlock(); + + } + } + }catch(IOException e) { + e.printStackTrace(); + }finally { + try { + lis.close(); + } catch (IOException e) { + e.printStackTrace(); + } + try { + kos.close(); + } catch (IOException e) { + e.printStackTrace(); + } + } + }); + t1.start(); + t2.start(); + try { + t1.join(); + t2.join(); + } catch (InterruptedException e) { + e.printStackTrace(); + } + close(); + } + + @Override + public void accept(KLALBRemoteLine from, KLALBPacket u) { + try { + //System.out.println(this+" "+u); + switch (u.getType()) { + case KLALBPacket.RST: + //reseted.countDown(); + if(!isListening()) + close0(false); + break; + case KLALBPacket.DATAT: + DATATPacket dtp = (DATATPacket) u; + if(isListening()) { + if(dtp.getNumber()==0) { + synchronized (backlogQueue) { + + //controller.getStreamPortBinder().checkIsConnected(new Pair); + InetSocketAddress is=new InetSocketAddress(from.getRemoteVaddr(), dtp.getSport()); + AtomicBoolean ab=new AtomicBoolean(true); + for (Iterator iterator = backlogQueue.iterator(); iterator.hasNext();) { + Object[] objects = (Object[]) iterator.next(); + if(is.equals(objects[0])) { + ab.set(false); + break; + } + } + if(ab.get()) { + if(controller.getStreamPortBinder().checkIsConnect(this,new InetSocketAddress(from.getRemoteVaddr(), dtp.getSport()))) { + ab.set(false); + } + } + + + if(ab.get()) + if(backlogQueue.offer(new Object[] { is,from,dtp})) { + /* controller.sendPacketToAddress(from.getRemoteVaddr(), new ACKTPacket(dtp.getDport(), dtp.getSport(),dtp.getNumber(),true,0), + 0,2);*/ + }else { + controller.sendPacketToAddress(from.getRemoteVaddr(), new RSTPacket(dtp.getDport(), dtp.getSport()), + 0,2); + } + + + } + }else { + controller.sendPacketToAddress(from.getRemoteVaddr(), new RSTPacket(dtp.getDport(), dtp.getSport()), + 0,2); + } + }else { + controller.sendPacketToAddress(from.getRemoteVaddr(), new ACKTPacket(dtp.getDport(), dtp.getSport(), dtp.getNumber(), + countInputBytes() < inputchachesize,dtp.getSendcount()), 0,2); + synchronized (inputchache) { + if (dtp.getNumber() >= inputcount) { + int currindex=(int) (dtp.getNumber()-inputcount); + int reqsize=1+currindex; + while(reqsize>inputchache.size()) { + inputchache.add(new AtomicInteger()); + } + for (int i = 0; i < currindex; i++) { + Object o=inputchache.get(i); + if(o instanceof AtomicInteger) { + ((AtomicInteger) o).incrementAndGet(); + if(((AtomicInteger) o).get()==50) { + controller.sendPacketToAddress(from.getRemoteVaddr(), new NACKTPacket(dtp.getDport(), dtp.getSport(), inputcount+i), 0,1); + System.out.println("请求快速重传:"+(inputcount+i)); + } + } + } + inputchache.set(currindex, dtp); + + Iteratoritr=inputchache.iterator(); + while (itr.hasNext()) { + Object datatPacket = itr.next(); + if(datatPacket instanceof DATATPacket) { + itr.remove(); + sendDeque.add((DATATPacket) datatPacket); + LockSupport.unpark(sendDequeLock); + inputcount++; + }else { + break; + } + + } + + /* int sze=inputchache.size(); + for (int i = 0; i < sze+1; i++) { + if(i < sze) { + DATATPacket klalbBlock = inputchache.get(i); + if(dtp.getNumber()itr=inputchache.iterator(); + while (itr.hasNext()) { + DATATPacket datatPacket = (DATATPacket) itr.next(); + if(datatPacket.getNumber()==inputcount) { + itr.remove(); + synchronized (sendDeque) { + sendDeque.add(datatPacket); + sendDeque.notifyAll(); + } + inputcount++; + inputcross=0; + }else { + inputcross++; + if(inputcross>5) { + inputcross=0; + controller.sendPacketToAddress(from.getRemoteVaddr(), new NACKTPacket(dtp.getDport(), dtp.getSport(), inputcount), 32768); + System.out.println("请求快速重传:"+inputcount); + } + break; + } + } + + */ + /* + if (dtp.getNumber() >= inputcount) { + inputchache.add(dtp); + while (true) { + DATATPacket kkb = null; + for (int i = 0; i < inputchache.size(); i++) { + DATATPacket klalbBlock = inputchache.get(i); + if (klalbBlock.getNumber() == inputcount) { + inputchache.remove(i); + i--; + kkb = klalbBlock; + break; + } + } + if (kkb == null) + break; + sendDeque.add(kkb); + synchronized (sendDeque) { + sendDeque.notifyAll(); + + } + inputcount++; + } + } + */ + } + } + //dbg.println(from.getMonitor()+","+dtp.getNumber()); + } + break; + case KLALBPacket.ACKT: + ACKTPacket ackt = (ACKTPacket) u; + if(isListening()) { + controller.sendPacketToAddress(from.getRemoteVaddr(), new RSTPacket(ackt.getDport(), ackt.getSport()), + 0,2); + }else { + avaliable = ackt.isAvaliable(); + DATATPacket kl=null; + synchronized (sendmap) { + + kl=sendmap.remove(ackt.getNumber()); + + } + /*synchronized (sendlist) { + int val=KLALBUtils.binarySearchDATATPacketNumber(sendlist, ackt.getNumber()); + System.out.println(val); + if(val!=-1) { + kl=sendlist.remove(val); + } + }*/ + /*synchronized (sendlist) { + int val=0; + for (Iterator iterator = sendlist.iterator(); iterator.hasNext();val++) { + DATATPacket inetSocketAddress = (DATATPacket) iterator.next(); + if(inetSocketAddress.getNumber()==ackt.getNumber()) { + iterator.remove(); + kl=inetSocketAddress; + break; + } + } + //System.out.println(val); + }*/ + if(kl!=null) { + if(sendthread!=null) + LockSupport.unpark(sendthread); + controller.removeFromSend(from.getRemoteVaddr(),kl); + DATATPacket.arrayRecycle.recycle(kl.getData()); + if(kl.getSendRecord().size()==1) { + long RTTC=ackt.getRcvtime()- kl.getSndtime(); + if(RTTC>RTTAvg) { + RTTAvg=(RTTAvg+RTTC)/2; + }else { + RTTAvg= (RTTAvg*99+RTTC)/100; + } + //System.out.println(RTTAvg); + } + } + } + break; + case KLALBPacket.NACKT: + NACKTPacket nackt=(NACKTPacket) u; + DATATPacket st=null; + synchronized (sendmap) { + st=sendmap.get(nackt.getNumber()); + } + /*for (Iterator iterator = sendlist.iterator(); iterator.hasNext();) { + DATATPacket st = (DATATPacket) iterator.next(); + if(st.getNumber()==nackt.getNumber()) { + + } + }*/ + if(st!=null) { + controller.sendPacketToAddress(remoteaddr,st, 5-1); + } + break; + default: + break; + } + } catch (IOException e) { + e.printStackTrace(); + } + } + + @Override + public InetAddress getRemoteInetAddress() { + return remoteaddr; + } + + @Override + public InetAddress getLocalInetAddress() { + return localaddr; + } + + @Override + public void setLocalPort(int i) { + localport=i; + } + } diff --git a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java index 2999198..b7f3dfb 100644 --- a/src/org/kne/cloud/network/klalb/LineDecitionComparator.java +++ b/src/org/kne/cloud/network/klalb/LineDecitionComparator.java @@ -1,5 +1,6 @@ package org.kne.cloud.network.klalb; +import java.net.SocketException; import java.util.Comparator; import java.util.HashMap; import java.util.Iterator; @@ -7,28 +8,23 @@ import java.util.List; import java.util.Map; import java.util.concurrent.atomic.AtomicLong; -public class LineDecitionComparator implements Comparator { - 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(); - 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()); - } - }); - al.addAndGet(curr.getLength()); +public class LineDecitionComparator implements Comparator { + private Map predictedlatencys=new HashMap<>(); + public LineDecitionComparator (List krs,KLALBPacket curr,int priority) { + for (Iterator iterator = krs.iterator(); iterator.hasNext();) { + KLALBRemoteLine klalbRemoteSocket = (KLALBRemoteLine) iterator.next(); + long x=klalbRemoteSocket.getSndDelayFactor(); + //long x=klalbRemoteSocket.getMonitor().getLatency()>>1; + /*long al=0; + al=klalbRemoteSocket.statLengthBefore(priority)+curr.getLength(); long speed=klalbRemoteSocket.getMonitor(). getOutSpeed(); if(speed==0) { - if(al.get()>0) { + if(al>0) { x=Long.MAX_VALUE; } }else { - x+=al.get()*1000000000.0/speed; - } + x+=al*1000000000L/speed; + }*/ /*System.out.println(klalbRemoteSocket); System.out.println(klalbRemoteSocket.getSendQueue().size()); @@ -36,14 +32,14 @@ public class LineDecitionComparator implements Comparator { predictedlatencys.put(klalbRemoteSocket, x); } } - public Map getPredictedlatencys() { + public Map getPredictedlatencys() { return predictedlatencys; } @Override - public int compare(KLALBRemoteSocket o1, KLALBRemoteSocket o2) { - double t1=predictedlatencys.get(o1); - double t2=predictedlatencys.get(o2); + public int compare(KLALBRemoteLine o1, KLALBRemoteLine o2) { + long t1=predictedlatencys.get(o1); + long t2=predictedlatencys.get(o2); if(t1>t2) { return 1; }else if(t1>1; + this.latencyAvg=(latencyAvg*9+ newlatency)/10; //} + } + + if(latencyMin==Long.MIN_VALUE) { + latencyMin=newlatency; + }else { + /*if(newlatencyinSpeedMax) { + /*if(inSpeed>inSpeedMax) { inSpeedMax=inSpeed; }else { inSpeedMax=(inSpeedMax*9999+inSpeed)/10000; @@ -128,11 +187,22 @@ public class Monitor { outSpeedMax=outSpeed; }else { outSpeedMax=(outSpeedMax*9999+outSpeed)/10000; - } + }*/ + + this.inSpeedAvg=(inSpeedAvg*9+ inSpeed)/10; + this.outSpeedAvg=(outSpeedAvg*9+ outSpeed)/10; if(changeListener!=null) { changeListener.accept(this); } + + } + public void updateOutSpeedMax() { + if(outSpeed>outSpeedMax) { + outSpeedMax=outSpeed; + }else { + outSpeedMax=(outSpeedMax*9999+outSpeed)/10000; + } } private ConsumerchangeListener; @@ -142,22 +212,26 @@ public class Monitor { public void setChangeListener(Consumer changeListener) { this.changeListener = changeListener; } - public long getLatency() { - return latency; + public long getLatencyAvg() { + return latencyAvg; } public String toString() { StringBuilder sb=new StringBuilder(); + sb.append(name); + sb.append('\t'); 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"); + 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(latencyAvg/1000000).append("ms"); return sb.toString(); } public String toString2() { StringBuilder sb=new StringBuilder(); + sb.append(name); + sb.append('\n'); 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"); + 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(latencyAvg/1000000).append("ms").append("\t").append(jitter/1000000).append("ms\t"); if(state==OFFLINE) { sb.append((coolingTime-(System.currentTimeMillis()-mls))/1000L); diff --git a/src/org/kne/cloud/network/klalb/MonitoredSocket.java b/src/org/kne/cloud/network/klalb/MonitoredSocket.java index 8cf78f4..92e8974 100644 --- a/src/org/kne/cloud/network/klalb/MonitoredSocket.java +++ b/src/org/kne/cloud/network/klalb/MonitoredSocket.java @@ -8,7 +8,7 @@ import java.util.Timer; import java.util.TimerTask; import java.util.concurrent.atomic.AtomicLong; -import org.kne.cloud.network.mport.FilterSocket; +import org.kne.cloud.network.FilterSocket; public class MonitoredSocket extends FilterSocket { public Monitor getMonitor() { diff --git a/src/org/kne/cloud/network/klalb/PONGPacket.java b/src/org/kne/cloud/network/klalb/PONGPacket.java index aa45750..be88e6e 100644 --- a/src/org/kne/cloud/network/klalb/PONGPacket.java +++ b/src/org/kne/cloud/network/klalb/PONGPacket.java @@ -5,39 +5,53 @@ import java.io.DataOutput; import java.io.IOException; public class PONGPacket extends KLALBPacket { - public long getTime() { - return time; + + public long getTimepingsnd() { + return timepingsnd; } - private long time; + + public long getTimepingrcv() { + return timepingrcv; + } + + public long getTimepongsnd() { + return timepongsnd; + } + private long timepingsnd,timepingrcv,timepongsnd; @Override protected void writeToStream(DataOutput dto) throws IOException { super.writeToStream(dto); - dto.writeLong(time); + dto.writeLong(timepingsnd); + dto.writeLong(timepingrcv); + dto.writeLong(timepongsnd); } @Override protected void readFromStream(DataInput din) throws IOException { super.readFromStream(din); - time=din.readLong(); + timepingsnd=din.readLong(); + timepingrcv=din.readLong(); + timepongsnd=din.readLong(); } @Override public long getLength() { - return super.getLength()+8; + return super.getLength()+24; } public PONGPacket() { super(PONG); } - public PONGPacket(long time) { + + public PONGPacket( long timepingsnd, long timepingrcv, long timepongsnd) { super(PONG); - this.time=time; + this.timepingsnd = timepingsnd; + this.timepingrcv = timepingrcv; + this.timepongsnd = timepongsnd; } + @Override public String toString() { return "PONG"; } - public void redeltaTime(long delta) { - time+=delta; - } } diff --git a/src/org/kne/cloud/network/klalb/RSTPacket.java b/src/org/kne/cloud/network/klalb/RSTPacket.java index 6a0ba15..2e42556 100644 --- a/src/org/kne/cloud/network/klalb/RSTPacket.java +++ b/src/org/kne/cloud/network/klalb/RSTPacket.java @@ -4,7 +4,7 @@ import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; -public class RSTPacket extends KLALBPacket { +public class RSTPacket extends KLALBPacket implements PortPacket{ private int sport,dport; public RSTPacket(int sport,int dport) { super(RST); diff --git a/src/org/kne/cloud/network/klalb/SendTask.java b/src/org/kne/cloud/network/klalb/SendTask.java index e738a96..ea54dc6 100644 --- a/src/org/kne/cloud/network/klalb/SendTask.java +++ b/src/org/kne/cloud/network/klalb/SendTask.java @@ -58,8 +58,7 @@ public class SendTask { public void check(int number) throws IOException { - long limit= (1<<(2*Math.min(4,count.get()-1)))*(number<3?200:1000*number)*1000000L; - //long limit=(1<{ - KLALBVirtualSocket kvs=null; - try { - s.setTcpNoDelay(true); - kvs=new KLALBVirtualSocket(kc, kr.getRemoteVaddr(), 23333); - 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 { - - try { - s.close(); - } catch (IOException e) { - e.printStackTrace(); - } - if(kvs!=null) - try { - kvs.close(); - } catch (IOException e) { - e.printStackTrace(); - } - } - }); - tl.open(); + new SocketToSocketProxy(new MultipurposeSocketAddress(ap.getProperty("local")), vmsa); System.out.println("提示:输入state并回车可以查看当前线路状态"); while(true) { String s=scn.nextLine(); switch(s) { case "state": 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(); - System.out.println(hostPort.getValue() .toString2()); + for (Iterator iterator = kc.getLines().iterator(); iterator.hasNext();) { + KLALBRemoteLine hostPort = iterator.next(); + System.out.println(hostPort .toString2()); //System.out.println(); } break; diff --git a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java index 8027082..df6a136 100644 --- a/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java +++ b/src/org/kne/cloud/network/klalb/SimpleKLALBServer.java @@ -4,34 +4,33 @@ import java.io.BufferedReader; import java.io.File; import java.io.FileReader; import java.io.IOException; -import java.net.InetAddress; -import java.net.ServerSocket; import java.net.Socket; import java.net.UnknownHostException; import java.util.Iterator; -import java.util.List; import java.util.Scanner; -import org.kne.cloud.network.klalb.LineManager.LineEntry; -import org.kne.cloud.network.mport.MultipurposeSocketAddress; -import org.kne.cloud.network.mport.Protocol; -import org.kne.cloud.network.mport.ProtocolDetectorServerSocketFactory; -import org.kne.cloud.network.mport.ProtocolDetectorSocket; -import org.kne.cloud.network.mport.ProxyProfileAnalyser; -import org.kne.cloud.network.mport.ProxyProfileExecutor; -import org.kne.cloud.network.mport.SocketBridge; -import org.kne.cloud.network.mport.TCPListener; +import org.kne.cloud.network.DatagramSocketListener; +import org.kne.cloud.network.MultipurposeSocketAddress; +import org.kne.cloud.network.ProtocolDetectorServerSocketFactory; +import org.kne.cloud.network.ProtocolDetectorSocket; +import org.kne.cloud.network.SocketBridge; +import org.kne.cloud.network.SocketListener; +import org.kne.cloud.network.SocketToSocketProxy; +import org.kne.cloud.network.SocketType; public class SimpleKLALBServer { - public static TCPListener tcpl,tcpl2; + public static SocketListener tcpl; + public static DatagramSocketListener udpl; public static ServerPropties sp; public static KLALBController kc; static { - MultipurposeSocketAddress.getServerSocketFactoryRegister().put("DETTCP", new ProtocolDetectorServerSocketFactory()); + MultipurposeSocketAddress.getSocketTypeRegister().put("DETTCP",new SocketType(null, new ProtocolDetectorServerSocketFactory())); } public static void main(String[] args) throws IOException { + //Debuger dbg=new Debuger(); + //dbg.start(); - System.out.println("KLALB负载均衡V2.0"); + System.out.println(CONST.klalb+" V"+CONST.klalbver); sp=new ServerPropties(); System.out.println("虚拟地址:"+sp.getVirtualIP().getHostAddress()); @@ -58,10 +57,11 @@ public class SimpleKLALBServer { System.out.println("reload:重新加载线路配置"); break; case "state": - LineManager le=kc.getLineManager(); - for (Iterator> iterator = le.entrySet().iterator(); iterator.hasNext();) { - java.util.Map.Entry hostPort = iterator.next(); - System.out.println(hostPort.getValue() .toString()); + System.out.println("状态\t上传流量\t下载流量\t上传速度\t下载速度\t延迟\t抖动\t下一次重试"); + for (Iterator iterator = kc.getLines().iterator(); iterator.hasNext();) { + KLALBRemoteLine hostPort = iterator.next(); + System.out.println(hostPort .toString2()); + //System.out.println(); } break; default: @@ -72,13 +72,21 @@ public class SimpleKLALBServer { private static void openPort(String bip) throws IOException { if(tcpl!=null) tcpl.close(); + if(udpl!=null) { + udpl.close(); + } MultipurposeSocketAddress mpsa=new MultipurposeSocketAddress(bip,"DETTCP"); - tcpl=new TCPListener(mpsa); + tcpl=new SocketListener(mpsa); tcpl.setCon((soc)->{ ProtocolDetectorSocket pds=(ProtocolDetectorSocket) soc; if(!pds.getProtocolStack().isEmpty()&&pds.getProtocolStack().pop().getName().equals("KLALB")) { - KLALBRemoteSocket krs=new KLALBRemoteSocket(pds); - kc.addRemoteSocket(krs); + KLALBRemoteLine krs=null; + try { + krs = new KLALBRemoteLine(new StreamKLALBPacketLink(pds)); + kc.addRemoteLine(krs); + } catch (IOException e) { + e.printStackTrace(); + } }else { MultipurposeSocketAddress mpsa2=new MultipurposeSocketAddress(sp.getLocal()); Socket s=null; @@ -108,49 +116,21 @@ public class SimpleKLALBServer { } } }); - tcpl.open(); - } - private static void openLocalPort(String bip) throws IOException { - if(tcpl2!=null) - tcpl2.close(); - MultipurposeSocketAddress mpsa=new MultipurposeSocketAddress(bip); - tcpl2=new TCPListener(new MultipurposeSocketAddress("KLALB", "::0", 23333)); - tcpl2.setCon((soc)->{ - Socket s=null; + udpl=new DatagramSocketListener(new MultipurposeSocketAddress(bip, "UDP")); + udpl.setCon((r)->{ + KLALBRemoteLine krl; try { - s=mpsa.connectSocket(); - 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(); + krl=new KLALBRemoteLine(new DatagramKLALBPacketLink(r)); + kc.addRemoteLine(krl); } catch (IOException e) { // TODO 自动生成的 catch 块 e.printStackTrace(); - }finally { - if(s!=null) - try { - s.close(); - } catch (IOException e1) { - // TODO 自动生成的 catch 块 - e1.printStackTrace(); - } - try { - soc.close(); - } catch (IOException e) { - // TODO 自动生成的 catch 块 - e.printStackTrace(); - } } - }); - tcpl2.open(); + } + private static void openLocalPort(String bip) throws IOException { + MultipurposeSocketAddress mpsa=new MultipurposeSocketAddress(bip); + new SocketToSocketProxy(new MultipurposeSocketAddress("KLALB_Stream", "::0", 23333),mpsa); } private static String fileRead(String filePath){ //1.定义一个BufferedReader对象,将文件内容读取到缓存 diff --git a/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java b/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java deleted file mode 100644 index 2da4c94..0000000 --- a/src/org/kne/cloud/network/mport/SocketToSocketProxyProfileEntry.java +++ /dev/null @@ -1,16 +0,0 @@ -package org.kne.cloud.network.mport; - -public class SocketToSocketProxyProfileEntry extends ProxyProfileEntry { - public SocketToSocketProxyProfileEntry(String l, String r) { - src=new MultipurposeSocketAddress(l); - des=new MultipurposeSocketAddress(r); - } - private MultipurposeSocketAddress src; - private MultipurposeSocketAddress des; - public MultipurposeSocketAddress getSrc() { - return src; - } - public MultipurposeSocketAddress getDes() { - return des; - } -} diff --git a/src/org/kne/cloud/network/nathole/NatholeTestS.java b/src/org/kne/cloud/network/nathole/NatholeTestS.java index e08c3bf..488b4c2 100644 --- a/src/org/kne/cloud/network/nathole/NatholeTestS.java +++ b/src/org/kne/cloud/network/nathole/NatholeTestS.java @@ -9,12 +9,12 @@ import java.nio.channels.SocketChannel; public class NatholeTestS { public static void main(String[] args) throws IOException { - for (int i1 = 1024; i1 < 65536; i1++) { + for (int i1 = 1024; i1 < 32768; 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.connect(new InetSocketAddress("183.199.49.116",21000)); sc.close(); int i2=i1; new Thread(()->{ diff --git a/src/org/kne/debug/TimeDebugger.java b/src/org/kne/debug/TimeDebugger.java index c9f3a79..c9631d8 100644 --- a/src/org/kne/debug/TimeDebugger.java +++ b/src/org/kne/debug/TimeDebugger.java @@ -3,15 +3,16 @@ package org.kne.debug; import java.util.ArrayList; import java.util.HashMap; import java.util.Hashtable; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; public class TimeDebugger { - Map l=new Hashtable(); + Map l=new LinkedHashMap(); long tmp=-1; public void putTime(String name){ - long v=System.currentTimeMillis(); + long v=System.nanoTime(); if(tmp==-1){ l.put(name, 0l); }else{ @@ -20,6 +21,6 @@ public class TimeDebugger { tmp=v; } public void print(){ - System.err.println(l+"(ms)"); + System.err.println(l+"(ns)"); } }