From 5c262f1d3c759a40934675396d0e36e7380eed01 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 8 Jun 2018 16:49:23 -0400 Subject: [PATCH] GH-66: Add DynamoDbLockRegistry implementation (#93) * GH-66: Add DynamoDbLockRegistry implementation Fixes https://github.com/spring-projects/spring-integration-aws/issues/66 * * Remove `Lifecycle` from `DynamoDbLockRegistry` in favor of a thread execution in the `afterPropertiesSet()` * Fix `lock()` interruptibility logic * * Remove `mavenLocal()` since the upstream PR is merged * decrease an amount of expectations in the Kinesis test * * Upgrade to SC-AWS-2.0.0.RC2 * * Add Docs for the `DynamoDbLockRegistry` --- README.md | 10 + build.gradle | 6 +- gradle/wrapper/gradle-wrapper.jar | Bin 54329 -> 54413 bytes gradle/wrapper/gradle-wrapper.properties | 2 +- .../aws/lock/DynamoDbLockRegistry.java | 538 ++++++++++++++++++ .../integration/aws/lock/package-info.java | 4 + ...nesisMessageDrivenChannelAdapterTests.java | 2 +- ...amoDbLockRegistryLeaderInitiatorTests.java | 249 ++++++++ .../aws/lock/DynamoDbLockRegistryTests.java | 319 +++++++++++ 9 files changed, 1126 insertions(+), 4 deletions(-) create mode 100644 src/main/java/org/springframework/integration/aws/lock/DynamoDbLockRegistry.java create mode 100644 src/main/java/org/springframework/integration/aws/lock/package-info.java create mode 100644 src/test/java/org/springframework/integration/aws/leader/DynamoDbLockRegistryLeaderInitiatorTests.java create mode 100644 src/test/java/org/springframework/integration/aws/lock/DynamoDbLockRegistryTests.java diff --git a/README.md b/README.md index 3de3bab..fc11f72 100644 --- a/README.md +++ b/README.md @@ -617,6 +617,15 @@ this.amazonKinesis = AmazonKinesisAsyncClientBuilder.standard() Where you should specify the port on which you have ran the Kinesalite service. Also you can use for you testing purpose a copy of `org.springframework.integration.aws.KinesisLocalRunning` in the `/test` directory of this project. +## Lock Registry for Amazon DynamoDB + +Starting with _version 2.0_, the `DynamoDbLockRegistry` implementation is available. +Certain components (for example aggregator and resequencer) use a lock obtained from a `LockRegistry` instance to ensure that only one thread is manipulating a group at a time. +The `DefaultLockRegistry` performs this function within a single component; you can now configure an external lock registry on these components. +When used with a shared `MessageGroupStore`, the `DynamoDbLockRegistry` can be use to provide this functionality across multiple application instances, such that only one instance can manipulate the group at a time. +This implementation can also be used for the distributed leader elections using a [LockRegistryLeaderInitiator][]. +The `com.amazonaws:dynamodb-lock-client` dependency must be present to make a `DynamoDbLockRegistry` working. + [Spring Cloud AWS]: https://github.com/spring-cloud/spring-cloud-aws [AWS SDK for Java]: http://aws.amazon.com/sdkforjava/ [Amazon Web Services]: http://aws.amazon.com/ @@ -632,3 +641,4 @@ Also you can use for you testing purpose a copy of `org.springframework.integrat [Kinesalite]: https://github.com/mhart/kinesalite [Amazon SQS Message Attributes]: https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-message-attributes.html [Amazon SNS Message Attributes]: https://docs.aws.amazon.com/sns/latest/dg/SNSMessageAttributes.html +[LockRegistryLeaderInitiator]: https://docs.spring.io/spring-integration/docs/current/reference/html/messaging-endpoints-chapter.html#leadership-event-handling diff --git a/build.gradle b/build.gradle index f940376..02a447d 100644 --- a/build.gradle +++ b/build.gradle @@ -32,10 +32,11 @@ repositories { ext { assertjVersion = '3.9.1' + dynamodbLockClientVersion = '1.0.0' servletApiVersion = '3.1.0' log4jVersion = '2.11.0' - springCloudAwsVersion = '2.0.0.M4' - springIntegrationVersion = '5.0.4.RELEASE' + springCloudAwsVersion = '2.0.0.RC2' + springIntegrationVersion = '5.0.6.BUILD-SNAPSHOT' idPrefix = 'aws' @@ -106,6 +107,7 @@ dependencies { compile('com.amazonaws:aws-java-sdk-kinesis', optional) compile('com.amazonaws:aws-java-sdk-dynamodb', optional) + compile("com.amazonaws:dynamodb-lock-client:$dynamodbLockClientVersion", optional) compile("javax.servlet:javax.servlet-api:$servletApiVersion", provided) diff --git a/gradle/wrapper/gradle-wrapper.jar b/gradle/wrapper/gradle-wrapper.jar index f6b961fd5a86aa5fbfe90f707c3138408be7c718..91ca28c8b802289c3a438766657a5e98f20eff03 100644 GIT binary patch delta 7397 zcmY+JWmFVEyv3IV=@jV(>F!z@2@w$KmS%wkX_s20k!ES6J0t|8k&cC>kxuDOX?#B4 zJ1^$b@7_Ce=FH5Ong2AGqQ;b=#*3n*Sr_*uNE)JFxShG70OBcY>n*sk%p!xW6fgjQ zBRDM&7tC8npX5nn+ckXXnY>Y{cCI=kWN6%){zd#LVaB+K7p4${BAbJYEe{+=^g7o2 z__WmJ>zD&~hr|ATWSj$-S%Gb#e5Sm?p<;&4WZ3+X?4hqH!1wr#Ksqkbh>`eCw*T+> zdr4oIU5+5{&z^U6jHJ3~r{01DW22Zt>hf#9IBE|C7t*QCvC|Q0*{0hOPzIC%= zQ~86R@Nx1*(OD6DAAhX2>zDp?^3l(m3<+(5VctWD-UDZ}ZTc@y9Q*FiP_%Axap|1< z!cW)r{Ltu)8r7|NIH$V8}$P&4CK|U7pkjTPh7{DVY2-*0E<=ux<`7XB~s7k>OnJOT0zW_N)}E z+g1`<4W42#8BC;7L%Z2UF=_s>L?z*|w|n&6dGuDALeY@xl%#dE8tsB$VoFb#-B1$; z?l2EYZZ3fYzSNvjt^GTfLZ5Ck_rS zb=+}lXfza2Kg4%j^mWuVtR-m0-PGJ-<{|A1_w>f(M2+>%gZeS4Q=RPJtsOa~wu^M{;luOUT~@`NyIiV`8gW1%!)1FBg69GkONB9Yy5ArTOk<9HJqOLnS; z_ha%bf9_T|HUSg07;>0*`)&S?nJd$ zX35$0rU2z3=pctjJUdOtFihi$tv5`)U{}Il=z(>9A?GY3_6KH&9M;c_J`G%m=(d@~ zQT|Z%r%J9eROrRbV?3?ln9u_v;qBZ?{YznP790&%Huf27@X93tC6JT;OaECsKUiea zzsl(|B)p_fa=uWSn^5<*HOAGAd>a^5NZyd3eI-Dg=$i<;SI>EGT@Vi9(jt~1bqU#R)(|KcV?Vnrdv z>od_H7eYVBMcSVUI!l{l9kdDQIF4F^RZM2ct(;m%MiI7=QXq_GefpEOfWhE~h@&Qi zRz{RPm*_pzQ9)$~R-*KwZ_rHOcl@xjB5G}jfbd-AAJN0GJ{*^1n+xVki0+BFneb56 z_OQA3_N?{?Hz7?YG`|5^|5vRjWbL!d@@)$E4_73(1!$ zexzL2VHAi_;ULOV67s}an3{HylUZb&MWy2J)T{QZ?vY-?fi9f|8RN-eyeop^y0-dz#>hdNrDDENnCbbIPam0LCGW!8rk8}OVsk1$ zeMXB>B6zg+;t#aLLHu1m6=DSY@hk^^P(jT7X>k(@nj-`AW#tb6^2l9Q32_x>)@d%N ze{vs?DtgQbR>Vh4e+Ut4j@2UYV`daonL6yRRonz7?d$hAB>k9>C(#fyKz=7vb#s?& zPG?|(#U;du3%CZEfRtIJHF(Tk5}UsyGiOkK8cqJw9cUEg?;v%o$G%i?=bA&=FK0Ab z%73UIXQ?6;in>pMTS3N`Z`RLn#euKic?WCsA~;jZYU@HfIlFrBRi8CJn`NeUig`p%kI zm6yJA34!)1L5d2pz+PDrXItk)oVs2~Z+f|kpFsvkmdi%^V4p4n`6EW96&6If%`F<;#pM$`8!{M*kv2ftii=> z)VSk{*&GKY?4ac_K_Ccj$4-qst|QIVS$F$}BArRSw%hKRJ!u^NsR<9(TkcfEm;e}2 zR7T#HY*O^0A!mmjWUj}BIc{Rc_ABb4&7ea~aj6I<;O!O2av>lyJLnBXsa|sjwk3|? z<$7m#H=r0{o23~ut7i#a$>+(jRsO!X75JvKt`y39dJ!7nocO5$MItV<$dGL9$}gdt z>V5ShA%*FMWseR0UV;WF4Eq13E>ak{0|v#w<*yH^ADvR0B--JYm8-5H1K zVCtf~$hfY>fx)s;iX+Lc6=IOnik?}#FSTyZ0Y046r_yQ{q5cUUYA`aXv#IbZ82a)o=5&_lnAI}NuapU)+> zvraAS9lK`~t!C;2wqJmXeNw#MXd-l4gqXa}R^QZnbTOJK;;?@I@rm%&)XT!SNUM{(Lnx-WhRxPu7{fnlcuStU#%j%GVhvSQyi=Oa>s@|Bti%&f;g2~wYZL4OB}Mxk z+b?F2QRYl8SEH$^Gq|gKB#haIz z!t;iN69p}snp1u63`$`ywSPlVpl8E$$3d3G`1YD}`;<3C3`sekPxy0hxMmFyo{Al` zjUy|SyiSr2y5f=5&mgP6XF%K}*ytW8IoHpLli1}$X7!wnQbBvW;f3szBEF}rR? zY`y}KF6S!8;D}f}))rrPH$rXQNJnS4R74kLLc$Xlfi3HZv7(c;kW~xC!l%tIi8<%%fuw4h?eLli`gI6Kc1dj;G?dZU%Rmk>0 zlI*fs%~LpWeDKVqjhU}l1BezIGOxA@6Fy1dK%$X2@wi`9-Dk2rw!iAbE@yVE;^ikoLwA7OmE8`D`Eh`+;V~q8G+M=nR8w0D3rUpl=QIqW6cR#W zY;H=fZDK3v(JfVx;PyNoNGh$v?A6D?Ny}xo_u*btdiuJqOFV-}gTehtE+HTLD#vCi z(P;=v&BQ)l*JOM5u$++%C0KV4T1u0cFKqo1W9t@;MmdYl7p_>lNlO)kT)JCW9VtPHV=LvN7Cmu!15P~= zF$c_e{H7EUclHaLPVW}W^}+i_m3k%^T=!F`!E7jtgD#J*Vrk%!gD!bR`;?-cQw6^s zMs?OLY~y(CPr6ho2v{&%;&-<-BtQw`Ip1zs{2 zlJHij73LF3dT9AFz7*)QEA_3p!}WYQ0?UXmjQacL??6~-aXabU#l!MXHFclTo1S(a{$9+Ou=$cZDW!05dm|L@K!HX_WTu+?pi z$kqV%mHXt^bQ5$ho@BSoFJ?IYgt;hefLyzkEix$*J}QA5*Dk9&DFAaCBd*-o7VNTW zlvG=i*DqY|*|1z;)URPCs(YsgPe~S0Zam|}Y9`eqd2n!l{gFqusPv4n+&W+Fk)xbm zdzXM7H=IUdvNTq4Kthztjw5VUA;BmZ#T_a} zXfCQ{AXqeAM&}`GrRdZ)&TYsdzqA9lM`6MdV_3Bc&Q<-{go7O~E^ zB+1H5?E^JZhaahb98T-Q=;L^R_V?-q9=(nTHYUTIvVQs?(oce8pj#_agXmxlPa_bK zau%s)QijfkTYKXL^IMmNlo8?0wTk{5qF^^Za9AiYC4CJlE5;ez-EDZhq=WtV+9|=J zdrjxAPf6+VIY|9?-N@&u?um#2njz=;MmR!S(Qmtp^<%eV{n8&8ahAq!gPbab^ic`s zyJ0~wcc99P7n5uLIQeEmYkzWh5Vrh@;jC-W2ZwsPj>On| z(^XgybF&i7-OamAiBIwSlZlSR#Vdp}R)ZsekYoaci@43^_HDzOFN#Ygfn8V?j~7T|6QKR zmkHh^#w|>z0HL?8tL4 z#j>Jk@asSr=>&I8HHPo-uHK!5P$s$*!T6DZaAxro&F<|k)!b3vsAMfhs+KaHbFTqF z@E*H_x@c_yC1XKvMRO`cO0DhJz=Ysa#sj0v%Cb_Lfw z__Tj&^5blO$_^Vs9$%Defz%#eJ>{M>`+`?KO3+@yNN&l@ph=^p?3kpN6HX^mLL8;h znX}4v-_c1ZGNxK)$1vi4km%f%ejMkj0D~6P=zkXprF=)BTwIrrvDK197V+Vd)W1H7wH&?}1G!1${&? zji^*OwX5nDnXEZsh`aiRBDWb5%DiG%R*%%4a7zk2KHuRt{-opP?bG6RQ@_Jrjwlcc z1c?I{V!mStTm$P+uMqeQJ+G&A-Bogkd9n^ zow6ZV5GhtRsc+n5Vpm!gL~@IgZ%M6SEirNsUCuhF?$hN(;FV7cQV-0WmV+mE&j1UM z*5gsoF|(+ckJqjHl>N;;_7d$6)awK+so|DP85m4m6cAC}YjgzQewH$pvS#KLNV=gt zWpcapyxLARfGB_63|p6Ui?{UbnZJMi1E7pVZdJRUG0y0E8+C|I2lvjlSp5_*-9ppR z9Q?Y|hjxxf;g3l!AMqNq?7hENyaue$&EQ)6FbDWW#-+fS((exFbIWL=*R@bj_n6DF zGCb@vMge`!wA!G4&-nEB;*{lqQ0^)M{2B->b;pmtJ_V>dbMl(ZYwlESZKFt3Bd@t; zzC*F~yF5cZKTtpDuTE;>X#^O@JfP2J!20&Rf0;m2)H9@oGAu5WhS2UQ705k;s>y`Sd{YI7DBC zBW86YCiE2dyyhHO54}mazQYIr&G>9TXvbTK9y3uYjZq3E0q|CoB)hdt{WX#^W6{9b z&nmv5v0eF#uJwL#X3rXVGUQea5`B~e}uKLq6{Vekx6e!b8jq)3LJg-M=&0@ZdMnC*XKnKQY#v0E4a#+rkFQG4PF7MDnu$0{d@2)I(fA!1#X({jHj zKz%G%jhilKF2nq6|EC<~QQ!(T#~2+bAc{|YD0lVv+2a^p7)&X3?Y#g)mTJl)9^m8s z*)qZWHG$;QO3&fF5Pr+Eyk+V;JgzHaiSU4Mr<=IkOW*N3%Xx*Cn8({WVH}|VjK0f+x=;||J9}N#t+mw ze-&XLYz)9#H4ID`K0E~8)~)&PQ>%`k2OI7G{A4i>C{zExYTMuX+yA!D!tkzMlE1|_ zJb1|f2XI2=pR7{Bg%1rA!qEmPf!1pOEH!mJxP}@z+-Q&k=&b(F${dseX6XFGlR;&m zkKsR5A5sMBnfy0mz^jMEfL1neP8()8K7s?K8!nCj0ncpU%{GGY_$fko3xX2G4?csi z2Qon5av{`k<6${qqw_y(>mc~oXMUIxNcQ2Mb?<@;ry1b@(t7=sOmL?Wejv8@KfLmy zgfIO~a(wy2|{?vC?!fAymk}} ztd0Gb$0wk}UyaEEixU1@lEOR3IDoB5|C3fG5&at~GERxo1lJmW56nvcpD;H4Kl>BH zffJNKuP^^>xd}-iGW#DkWz+s0$^pd5|0@~cc$0jXlaLj2=bn>46Jlg*Nq_8n4 delta 7284 zcmZ9Rbxa&UyY*p_;_gLisRLH+^J!mIhcBl0_*qd%+} z)z^`2t@btAW%fOwIG?H%1_Gy}x}Ohmufbm)bnYro1zK{}9Me&D!8f@=>j4?J0qgJA zg}{&N4ZX;wm8$Z$}t~&p!qz zM5;2nN)8r$+-$is#4E@!h2DJts4|?vM8xO@h=*PjKYn;=j4y+ja!IgSxH&u(A7W1;N33(Q2q!A~a znNEU3e@F;(X)!lnU|J{3)s6K=skjw?AmpuVeqrer_x^Qs?U&(N%rWdSjeR~{n!e@(0V97}K*+nu2Gq5=j6o^hv|lm0T9Bb8WMS zgKYy1bsTHu6|8vQ?<5DpPa|m~{O6AS^c_47VPgZ=b4afD;E_7_Xb%RrlN{zQ?;(Li zk{4cli|tG6EF)i*N7s||M~56IUUx4Wa+^=RB_F3rZ&~^KuljFRS0H}eL+JeqEg{}H zg-#~J!arLXb*9*e!eH!7QguKjMMk{*#Re%DD+RoOV+Vv3c? z%Hl`dk8~MmbT;yavPj;fG|3rJJ+TeqjV+kDyskGI#1j|bu@EVctW@VIiwMv498cGt z3~9Jq|71$c-rJ`=M1*Kk?gj!{PPI+l6LXxK(Zko3FVB}enFMsJ1scd_yis6sE&zQV z59eLQhBr$;2q zENT4}bGlE%Z?&PZn(k*cvMTfMGPUc5GPg3dMG`Tb0Zo2bZAMsZT_(XqJ>j;N*gs~4 zR>e9`XjM2!Z6K2Yxgoa9`a?fa&&8sA7C>E^p#7RnEM_e?n>+Wsc2$`Mvp~cu+x_GE z7~xu>RZ*PooqG`+#i9C*yjbD{(Tgkg2!SM$=r+WTC`k>s69#Eb^wiktBAd-Qh;air zLF@d4xu)Ou>3pZ0NWd;T%63rQ$E7_U_E0B8J)=a)1Mz8322(^MxmkEDO&G4fPmVRh zTB7+LNZ}xKL#M2AOk3b?IY$o$QAd@J7Vm9S$zmCqrLC6^tv9`W{fhGvVfLBo&B8ll z1#Vwm%#Fh8-?(1zQ45zOS{$2=eyEAawONDQ4`4WKE|t{uoO8~C%{xXt-q0{hEU@Ml zNMAHY4sxe$ra?)cEw!y9c|5|DMP;_#=9;OLz&*_DcUp05!Dk$KI!v;`#GtM)+;x*& zH1Z)5*(AA9Aib*i^aAfH;d|jk*M=bO{uk_#uzSQ`Z@Y5_-(;elL9~XZdG8)M)D#gA z3E<#RP~hO;B;oSVJ0$4w;oxq7aB$Rr9T#^SXbm3?gmCAxIOFZ^)Z}C%(uek1I=oRT zZ(b&{6zV1#YBMhCO?sb{_*I6fJVm@C?YG+c$HiW4gg%vJS|B`L=2ovCKlSTc7Z+_D z0*`9~8(a?j8(e^=5Nl^==M>4W+Zis_y*st99=(@=;A_T4f zT=>NdWLvRhY_A`22k9OX&n+qq*rz!@D zNbOF8e6^EoqI|I&&EqiihS_L#%$FUp zHH5}E+Yx9SDWq{yyagHZuXEHK&^>qDTTfERdpcCH46v7BktjXLcO zEgF3+W0D$Zxp6jfQKLRnF8Ma!+sTcIO@m87O}y-OXLWH53tCQ zH{cQxfsdM`o0Hu#8|VbiX&()wUuZWF-u>!ihJpn!7h_h%IcYdO z1{mQ3T*WVT&dT|1tT)3LYfQ8`eYtZ1?&vn65$iDAV9AR9=oCMKS;Jn<2$f9%8m>2X zdTY)PQ(k{!BNyPU{|ReyGrCN;{!9+}!AZz#RD3#DQ{va%(mB*R;k$=c?&cNaJGZ!X zV=bB;%C5tOR2(yQjVGV1(#t=B*R4##m~9a`dJierXH{O(Jeb?t<{e|;m(H&RD4cP#>8G zk4lf(3s8vck@u|0b~#DsCYvrq5P;^pH8~NE^8T9B)6KI38RZi z&4Vu<@bl`AJ!m~vI+IHtkf`XwycZm=JgN(~{GkQpIn?OSmxYd~B;i^hPW`Y7b&)|v zoQfS5)YY$YCK2~n_@<2LfX6jMx)-#-V9q6)z@*tkYgRYu`(4SIYilW^@N+|{CkufR zeU6u27BN?G^RebJPyAnYB)1jI@_M;SL|&}a&;3EfNQ zgMG>KWF9*(qMrOSVb@5 z6={LFVXask#>{ck;vqwfIXWJG#MmBs?+mqbVO=dHR?!O=&_!*cHo=2cneOa-9fo<3 zU%s*00RrT~KkdOzqO4$^1LT*TuYN%G@II}+oTKszMYMX3suww4wYoo5rtn;*Q>ARp zv#!5Ot2g~i<%L(b+y=+!NQPE#zKhyDO8`Rao&^h^OA*dv^_a<4gH(hgChjDggFTCS zDN{hnn83LJ&i7^fW^+R6eWVB%?R%nLt(}!S?=;Bi@At|2-44hRryBP;Pi$$YyFJY7 z1Gfqu1!K2y$k>%n)IDKFroQKqpjCbg*>?s!MgjFDN$EyRwwH0x%g^%t^7iN9bwy7w zFP{2?b1|is#^rbPIa+p@2YEUz`1)^}+pVw6p$Rv3=scGX)oK>*{tKYb6)(FZE;VVOs6uz?KPbMtFx9tVzX4 zEK01{Tn!42gf>7)zK6MA!D87dH;jKE+M$)|C>_(_i0in%db3?K$$)F}WX_2f0e-qeNH7EQ7=tIh-L}XuaF?; zDZHQy9hTD)`<4iNY9CXL-okbhZ#~?A={e~66>O`dz1gQs!@h(vi!-={osm8(B|#|c zdA1Hv2<}8at=;P)O`$Ade*e{RY}AApi?e=Frl=JZd5e-}qb1fq4SWhrgHW5Yl94;u z#p=hIA(%;h&RrBWoUo%rt&xT9KNj^**E|7E-^Tvr(xRApqI!hAFTaR;N|C>A=52^? zR4X6Tyw~~_W-j?e({%N72A2FidkoSr4qGlTZsbb)&Y&?Hjej`0#S~rV7q~D34dqc* z>BED%3e!&gaphbp8F7jkCfG%JnZs?V!i|ymNMEyAc?^2N{Ze$6sP%&SrRqZUo-O{m z-KucR<#HMn6tymcd8qeTdG-FKqyLNAy{O~}$Ne*nVq(=S=!Gy~;tSEA2@1$@ycb<; zo82-{$Y1{HETkRDVCL-seuOyaULfGp*q5D^WP`*W#=;xQ!{j8n>$d%E)lsQ@uhonl zo`o@ul5}4FF^?BiW0XSsL0k{|=NN{v5*NhsPuZzn!}%JU;OfPM3ex!mXb1eg(hpRm zl0k13EYj;pmRX!#D`z5a$S$PWru4gtuF$C;YhlGS_vidK_-~e%&{s6S-MMvjV)cIv zPiWtA70_W4+|p#{b~TNBNC-Oj-v2VQWoqn~o0Bw3b%1_=k{E63zEV9WYnytQu-}?0 zzVB*A7F{ZHmXSJ!Ia->WHQK>2(T493$)Sj*NG5tt^wq6w{=|J#!Z!wzx8o^PQtJHY zY2sO7sLAU(n+k?6K@ z!e%G8h>!0=7%r3V8&wBIE{!Jt6#dvWYqFBAB(2#>bI$7y{Uc=sC;reVV_ER}Ln?ct z<&R~eb7qM_Z?SuW<<=i+IR3lCm60(R-K4Tdv9fZVlFSEj*HqA38mQq$$S(+rYEBo=nYPVmu zdE~*fYRvN|YL?{?&eSSg$^x4JC`X_Bb>kbE?#6a-ximH&Bibi3s$~xmFjFOU{4X!c z5tItqJ!jT~&h^V!Cc_k`(Bxa5h3;7L94Aug8 zI)503EUDZw^Naj}408$aeP8A$aYOfF?N&N3W*~IsrIDRHmb4bwCpaw8pdBmYG*Z#Q znNPbLE2^oKm|FDZND6Dc9K>A)Z+6G;Sx&4;nS6M(3NXS%3liB&!&Ea~rjF1@zhIml zM<}H#cC*@vk_XlZN{2L$N%6iE4s7_r=7GW9?ArFc5h@W7j7wX7!mVT8PIx*icI5}O zz_f2*{H!G~ewBN!K=OFhuZGzlvvhM5&AA4JT zF3W*zCn_FfXO8p=vy&9`0T%8Y5*Vm!Alo)aIpc_Y)eSxAukGN4_Qtfq1)|qH>w*Oq ze|H}q4t&~2H{1RhY8M6U_D(8qOUIsxw_fLdEseSIY<$-?Q_zxOtsUv{Xt&D^DSACA zfz+eyJ)jkB*9OXoSMOIW^Ub)c9KqS~#k%MMEu{fWlGBpg!KR3WWk;X@_{$H zgR+mlL|qJw4I({z9;P?a8eO*j!Mf-zaS2ZgVy5CBx6cm00YtAU;M=kkHes@NC+*I7 zXR8wTslDI#v=0-a^<*c7g$VDPAJRiR_cb~ZXL#UXnHK$6Ouk&&>*#blvtU^2Ny#IS zNdaOr;m0NlM@M7ccqr+IqM=l11WUXbs6|RGb+L$e%UgWK8TiK92(|vr8PQh{hQN&8 zHv(4VcesAMn3)w^v}z9QVMZP~ECr?WBx2e8@|Ona3QyB&b~O#fJDl)qJ93=*A)saf zQ9~iWrCWNf9W^qEURF3gTWHcUvaJuim?#AFh8zQ-N2tWE1p&kR7gj&5kZwmLRmq6t zDe6^q0%7SMj$gCaVMPc`W%?_i|8voWOf^dlN#P+Gq^vNgFAqi%tlM5@n!Fip{A=z| zZe%lkadj+xQDTXsdSP2kRuNHE@j1$F*>z&dE7u)~#8~D&dIT8#jkEaN69s}b*Z_HR zi%16ErWZe6f^ILobOzyc!Z%|cH6koHD&9M?cax3{CvB43ij($u6Pn_-!6K!pgv^XjtvU~E$ACJ zx(81gk57hqsyO^6t#-O5tPMRvJNzwp*U)PfOt$*eN_LM~|Nd4*m=INkv=qX}scq2* ze!dcNF8Z1K&00xVs|9lhGF%cjk`MNlaUOOPBcV(rp`; z*pXNr_ou(W7C01FewCk)?7LebGPBcL>o=NKhB;GU2lV&(y7d#NHo(DwlMa< zuf)w~H{Ayua|;|-5dd{B_U~t`M!a=zgj~_+Vd4V@-Di78agKSG-|bP>vo*oyIIFCF zpNyk29Bdh%qkjcwyd7SBJFhz9^9}4jA96E7It@`0_Yx7$`qh7Gln@DVnOnhvd2DmD z^td?>Wqb|t>q{8$Jhv#l)ilrqOZtn;xKkAxzQV7EDD4vpiX^~^PtSr1oygSa$dtT* zrQjOg{9j)E%yQ1-^8r-hsiBiQJ1Qce4YS(ox?HUg<=?5?Z;T0Bh~P~SwlgrdRmQF< zmyV&jeG3r-AKK4f@h~7)+`2uXroM#l{4s%lyy%~FY;g&Qi_(LQxU75RzJxDr;Zmi{ z1nPQ_O}w(mZSRnJuvM76-yKV?dAX|pfvkduD1fOE^~WcVy}6TKG$$orx!bUxb6AD= zWcGcfJa1J}B8lfgg`nEQA~gft`N<-#1%l@RoE{Tgf6#$KBmx4&5`9(k(KZ+T38X93 ze7Zl`@0sC=oY&@h_D#|jMoFHzL+`Z$j!eQXWADV`oAoMb{pQ2lwlvc@XCL|>=i}7ps@ZPpD8dj;6XJ(M&K*eW#d|q{Fe}?HjA8R zHtjbOZEs9pNLnKxb=*WfF)uVd5vz#VfUo4|J5%WJs}MM1$h+}Y^QXk$JJ!obvy7&| z{B1d}SwrKP>SC{~S9UTJ^@9Zh(Z_@r)W&@okqX%ikJr$f^`aXSwS^~Qw5~57_FtzY z0Mb0PL{fdv3DHAp=X5q1Ia`0$#Nf)KUzSH}NpTsG52gztC$USxLdA_5+6i@Rf`+Mx z6rns6k=gGfUic=}iXe(B(jn2!b=KabIIrl>a>BaR=c#%fDaG9%Yfca|c!^#aJt?~c zUCrZO(GWV~9Oe$Fe-qIDs#$ZzwuGUz>xWFe&b<1~p1&3RDEul1QOp60*3x}sMBypn zzM;G--A?C_?ZYUeX6ktK!}0dMy{4w5>}XKM=vVo3Fh$|QUvM(j)x3&Y>x%W*9pCp0 zmS~1CW~R;-<#BMSf;=VJ?*Vb%s;b1hdxdYAmFM$ALx59Le1mczmU;=r8BqP#f@-lC zKtudrZ71&$Ig0-q#r;j={yWI|pc})4(324XK%o3T6#2#vz)}2%1 + * Can create table in DynamoDB if an external {@link AmazonDynamoDBLockClient} is not provided. + * + * @author Artem Bilan + * + * @since 2.0 + */ +public class DynamoDbLockRegistry implements ExpirableLockRegistry, InitializingBean, DisposableBean { + + /** + * The {@value DEFAULT_TABLE_NAME} default name for the locks table in the DynamoDB. + */ + public static final String DEFAULT_TABLE_NAME = "SpringIntegrationLockRegistry"; + + /** + * The {@value DEFAULT_PARTITION_KEY_NAME} default name for the partition key in the table. + */ + public static final String DEFAULT_PARTITION_KEY_NAME = "lockKey"; + + /** + * The {@value DEFAULT_SORT_KEY_NAME} default name for the sort key in the table. + */ + public static final String DEFAULT_SORT_KEY_NAME = "sortKey"; + + /** + * The {@value DEFAULT_SORT_KEY} default value for the sort key in the table. + */ + public static final String DEFAULT_SORT_KEY = "SpringIntegrationLocks"; + + /** + * The {@value DEFAULT_REFRESH_PERIOD_MS} default period in milliseconds between DB polling requests. + */ + public static final long DEFAULT_REFRESH_PERIOD_MS = 1000L; + + private static final Log logger = LogFactory.getLog(DynamoDbLockRegistry.class); + + private final Map locks = new ConcurrentHashMap<>(); + + private final CountDownLatch createTableLatch = new CountDownLatch(1); + + private final AtomicBoolean running = new AtomicBoolean(); + + private final AmazonDynamoDB dynamoDB; + + private final String tableName; + + private AmazonDynamoDBLockClient dynamoDBLockClient; + + private boolean dynamoDBLockClientExplicitlySet; + + private long readCapacity = 1L; + + private long writeCapacity = 1L; + + private String partitionKey = DEFAULT_PARTITION_KEY_NAME; + + private String sortKeyName = DEFAULT_SORT_KEY_NAME; + + private String sortKey = DEFAULT_SORT_KEY; + + private long refreshPeriod = DEFAULT_REFRESH_PERIOD_MS; + + private long leaseDuration = 20L; + + private long heartbeatPeriod = 5L; + + /** + * An {@link ExecutorService} to call {@link AmazonDynamoDBLockClient#releaseLock(LockItem)} + * in the separate thread when the current one is interrupted. + */ + private Executor executor = + Executors.newCachedThreadPool(new CustomizableThreadFactory("dynamodb-lock-registry-")); + + /** + * Flag to denote whether the {@link ExecutorService} was provided via the setter and + * thus should not be shutdown when {@link #destroy()} is called. + */ + private boolean executorExplicitlySet; + + + public DynamoDbLockRegistry(AmazonDynamoDB dynamoDB) { + this(dynamoDB, DEFAULT_TABLE_NAME); + } + + public DynamoDbLockRegistry(AmazonDynamoDB dynamoDB, String tableName) { + Assert.notNull(dynamoDB, "'dynamoDB' must not be null"); + Assert.hasText(tableName, "'tableName' must not be empty"); + + this.dynamoDB = dynamoDB; + this.tableName = tableName; + } + + public DynamoDbLockRegistry(AmazonDynamoDBLockClient dynamoDBLockClient) { + Assert.notNull(dynamoDBLockClient, "'dynamoDBLockClient' must not be null"); + + this.dynamoDBLockClient = dynamoDBLockClient; + this.dynamoDBLockClientExplicitlySet = true; + this.dynamoDB = null; + this.tableName = null; + } + + public void setReadCapacity(long readCapacity) { + this.readCapacity = readCapacity; + } + + public void setWriteCapacity(long writeCapacity) { + this.writeCapacity = writeCapacity; + } + + public void setPartitionKey(String partitionKey) { + Assert.hasText(partitionKey, "'partitionKey' must not be empty"); + this.partitionKey = partitionKey; + } + + /** + * Specify a name of the table attribute which is used as a sort key. + * @param sortKeyName the sort key attribute name to use. + */ + public void setSortKeyName(String sortKeyName) { + this.sortKeyName = sortKeyName; + } + + /** + * Specify a value for the sort key attribute of the lock item. + * @param sortKey the sort key value to use. + */ + public void setSortKey(String sortKey) { + this.sortKey = sortKey; + } + + public void setLeaseDuration(long leaseDuration) { + this.leaseDuration = leaseDuration; + } + + public void setHeartbeatPeriod(long heartbeatPeriod) { + this.heartbeatPeriod = heartbeatPeriod; + } + + public void setRefreshPeriod(long refreshPeriod) { + this.refreshPeriod = refreshPeriod; + } + + /** + * Set the {@link Executor}, where is not provided then a default of + * cached thread pool Executor will be used. + * @param executor the executor service + */ + public void setExecutor(Executor executor) { + this.executor = executor; + this.executorExplicitlySet = true; + } + + @Override + public void afterPropertiesSet() { + if (!this.dynamoDBLockClientExplicitlySet) { + AmazonDynamoDBLockClientOptions dynamoDBLockClientOptions = + AmazonDynamoDBLockClientOptions + .builder(this.dynamoDB, this.tableName) + .withPartitionKeyName(this.partitionKey) + .withSortKeyName(this.sortKeyName) + .withHeartbeatPeriod(this.heartbeatPeriod) + .withLeaseDuration(this.leaseDuration) + .build(); + + this.dynamoDBLockClient = new AmazonDynamoDBLockClient(dynamoDBLockClientOptions); + } + + this.leaseDuration = + (long) new DirectFieldAccessor(this.dynamoDBLockClient) + .getPropertyValue("leaseDurationInMilliseconds"); + + + this.executor.execute(() -> { + try { + if (!this.dynamoDBLockClientExplicitlySet) { + try { + this.dynamoDBLockClient.assertLockTableExists(); + return; + } + catch (LockTableDoesNotExistException e) { + if (logger.isInfoEnabled()) { + logger.info("No table '" + this.tableName + "'. Creating one..."); + } + } + + CreateDynamoDBTableOptions createDynamoDBTableOptions = + CreateDynamoDBTableOptions + .builder(this.dynamoDB, + new ProvisionedThroughput(this.readCapacity, this.writeCapacity), + this.tableName) + .withPartitionKeyName(this.partitionKey) + .withSortKeyName(this.sortKeyName) + .build(); + + AmazonDynamoDBLockClient.createLockTableInDynamoDB(createDynamoDBTableOptions); + } + + int i = 0; + // We need up to one minute to wait until table is created on AWS. + while (i++ < 60) { + if (this.dynamoDBLockClient.lockTableExists()) { + return; + } + else { + try { + // This is allowed minimum for constant AWS requests. + Thread.sleep(1000); + } + catch (InterruptedException e) { + ReflectionUtils.rethrowRuntimeException(e); + } + } + } + + logger.error("Cannot describe DynamoDb table: " + this.tableName); + } + finally { + // Release create table barrier either way. + // If there is an error during creation/description, + // we deffer the actual ResourceNotFoundException to the end-user active calls. + this.createTableLatch.countDown(); + } + }); + } + + private void awaitForActive() { + IllegalStateException illegalStateException = + new IllegalStateException( + "The DynamoDb table " + this.tableName + " has not been created during " + 60 + " seconds"); + try { + if (!this.createTableLatch.await(60, TimeUnit.SECONDS)) { + throw illegalStateException; + } + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw illegalStateException; + } + } + + @Override + public void destroy() throws Exception { + if (!this.executorExplicitlySet) { + ((ExecutorService) this.executor).shutdown(); + } + + if (!this.dynamoDBLockClientExplicitlySet) { + this.dynamoDBLockClient.close(); + } + } + + @Override + public Lock obtain(Object lockKey) { + Assert.isInstanceOf(String.class, lockKey, "'lockKey' must of String type"); + return this.locks.computeIfAbsent((String) lockKey, DynamoDbLock::new); + } + + @Override + public void expireUnusedOlderThan(long age) { + Iterator> iterator = this.locks.entrySet().iterator(); + long now = System.currentTimeMillis(); + while (iterator.hasNext()) { + Map.Entry entry = iterator.next(); + DynamoDbLock lock = entry.getValue(); + if (now - lock.lastUsed > age && !lock.delegate.isHeldByCurrentThread()) { + iterator.remove(); + } + } + } + + private final class DynamoDbLock implements Lock { + + private final ReentrantLock delegate = new ReentrantLock(); + + private final String key; + + // It is safe to use a shared instance - access is guaranteed by the delegate lock. + private final AcquireLockOptions.AcquireLockOptionsBuilder acquireLockOptionsBuilder; + + private LockItem lockItem; + + private volatile long lastUsed = System.currentTimeMillis(); + + private DynamoDbLock(String key) { + this.key = key; + this.acquireLockOptionsBuilder = + AcquireLockOptions.builder(this.key) + .withReplaceData(false) + .withSortKey(DynamoDbLockRegistry.this.sortKey) + .withTimeUnit(TimeUnit.MILLISECONDS) + .withRefreshPeriod(DynamoDbLockRegistry.this.refreshPeriod); + } + + private void rethrowAsLockException(Exception e) { + throw new CannotAcquireLockException("Failed to lock at " + this.key, e); + } + + @Override + public void lock() { + awaitForActive(); + + this.delegate.lock(); + + this.acquireLockOptionsBuilder + .withAdditionalTimeToWaitForLock(Long.MAX_VALUE - DynamoDbLockRegistry.this.leaseDuration); + + boolean wasInterruptedWhileUninterruptible = false; + + try { + while (true) { + try { + while (!doLock()) { + Thread.sleep(100); //NOSONAR + } + break; + } + catch (InterruptedException e) { + /* + * This method must be uninterruptible so catch and ignore + * interrupts and only break out of the while loop when + * we get the lock. + */ + wasInterruptedWhileUninterruptible = true; + } + catch (Exception e) { + this.delegate.unlock(); + rethrowAsLockException(e); + } + } + } + finally { + if (wasInterruptedWhileUninterruptible) { + Thread.currentThread().interrupt(); + } + } + + } + + @Override + public void lockInterruptibly() throws InterruptedException { + awaitForActive(); + + this.delegate.lockInterruptibly(); + + this.acquireLockOptionsBuilder + .withAdditionalTimeToWaitForLock(Long.MAX_VALUE - DynamoDbLockRegistry.this.leaseDuration); + + try { + while (!doLock()) { + Thread.sleep(100); //NOSONAR + if (Thread.currentThread().isInterrupted()) { + throw new InterruptedException(); + } + } + } + catch (InterruptedException ie) { + this.delegate.unlock(); + Thread.currentThread().interrupt(); + throw ie; + } + catch (Exception e) { + this.delegate.unlock(); + rethrowAsLockException(e); + } + } + + @Override + public boolean tryLock() { + awaitForActive(); + + try { + return tryLock(0, TimeUnit.MILLISECONDS); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return false; + } + } + + @Override + public boolean tryLock(long time, TimeUnit unit) throws InterruptedException { + awaitForActive(); + + long start = System.currentTimeMillis(); + + if (!this.delegate.tryLock(time, unit)) { + return false; + } + + this.acquireLockOptionsBuilder + .withAdditionalTimeToWaitForLock( + System.currentTimeMillis() - start + TimeUnit.MILLISECONDS.convert(time, unit)); + + boolean acquired = false; + try { + acquired = doLock(); + + if (!acquired) { + this.delegate.unlock(); + } + else { + this.lastUsed = System.currentTimeMillis(); + } + } + catch (Exception e) { + this.delegate.unlock(); + rethrowAsLockException(e); + } + + return acquired; + } + + private boolean doLock() throws InterruptedException { + boolean acquired; + if (this.lockItem != null) { + this.lockItem.sendHeartBeat(); + acquired = true; + } + else { + this.lockItem = + DynamoDbLockRegistry.this.dynamoDBLockClient + .tryAcquireLock(this.acquireLockOptionsBuilder.build()) + .orElse(null); + + acquired = this.lockItem != null; + } + + if (acquired) { + this.lastUsed = System.currentTimeMillis(); + } + + return acquired; + } + + @Override + public void unlock() { + if (!this.delegate.isHeldByCurrentThread()) { + throw new IllegalMonitorStateException("You do not own lock at " + this.key); + } + if (this.delegate.getHoldCount() > 1) { + this.delegate.unlock(); + return; + } + try { + if (Thread.currentThread().isInterrupted()) { + LockItem lockItemToRelease = this.lockItem; + DynamoDbLockRegistry.this.executor.execute(() -> + DynamoDbLockRegistry.this.dynamoDBLockClient.releaseLock(lockItemToRelease) + ); + } + else { + DynamoDbLockRegistry.this.dynamoDBLockClient.releaseLock(this.lockItem); + } + } + catch (Exception e) { + throw new DataAccessResourceFailureException("Failed to release lock at " + this.key, e); + } + finally { + this.lockItem = null; + this.delegate.unlock(); + } + } + + @Override + public Condition newCondition() { + throw new UnsupportedOperationException("DynamoDb locks don't support conditions."); + } + + @Override + public String toString() { + SimpleDateFormat dateFormat = new SimpleDateFormat("YYYY-MM-dd@HH:mm:ss.SSS"); + return "DynamoDbLock [lockKey=" + this.key + + ",lockedAt=" + dateFormat.format(new Date(this.lastUsed)) + + ", lockItem=" + this.lockItem + + "]"; + } + + } + +} diff --git a/src/main/java/org/springframework/integration/aws/lock/package-info.java b/src/main/java/org/springframework/integration/aws/lock/package-info.java new file mode 100644 index 0000000..34693f1 --- /dev/null +++ b/src/main/java/org/springframework/integration/aws/lock/package-info.java @@ -0,0 +1,4 @@ +/** + * Provides classes supporting lock registry. + */ +package org.springframework.integration.aws.lock; diff --git a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java index 44ef72d..511e7d1 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/KinesisMessageDrivenChannelAdapterTests.java @@ -207,7 +207,7 @@ public class KinesisMessageDrivenChannelAdapterTests { assertThat(n).isLessThan(100); // When resharding happens the describeStream() is performed again - verify(this.amazonKinesisForResharding, atLeast(2)).describeStream(any(DescribeStreamRequest.class)); + verify(this.amazonKinesisForResharding, atLeast(1)).describeStream(any(DescribeStreamRequest.class)); this.reshardingChannelAdapter.stop(); } diff --git a/src/test/java/org/springframework/integration/aws/leader/DynamoDbLockRegistryLeaderInitiatorTests.java b/src/test/java/org/springframework/integration/aws/leader/DynamoDbLockRegistryLeaderInitiatorTests.java new file mode 100644 index 0000000..dd6b5ed --- /dev/null +++ b/src/test/java/org/springframework/integration/aws/leader/DynamoDbLockRegistryLeaderInitiatorTests.java @@ -0,0 +1,249 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.aws.leader; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import org.junit.AfterClass; +import org.junit.BeforeClass; +import org.junit.ClassRule; +import org.junit.Test; + +import org.springframework.integration.aws.DynamoDbLocalRunning; +import org.springframework.integration.aws.lock.DynamoDbLockRegistry; +import org.springframework.integration.leader.Context; +import org.springframework.integration.leader.DefaultCandidate; +import org.springframework.integration.leader.event.LeaderEventPublisher; +import org.springframework.integration.support.leader.LockRegistryLeaderInitiator; +import org.springframework.scheduling.concurrent.CustomizableThreadFactory; + +import com.amazonaws.services.dynamodbv2.AmazonDynamoDBAsync; +import com.amazonaws.services.dynamodbv2.model.DescribeTableRequest; +import com.amazonaws.waiters.FixedDelayStrategy; +import com.amazonaws.waiters.MaxAttemptsRetryStrategy; +import com.amazonaws.waiters.PollingStrategy; +import com.amazonaws.waiters.Waiter; +import com.amazonaws.waiters.WaiterParameters; + +/** + * @author Artem Bilan + * + * @since 2.0 + */ +public class DynamoDbLockRegistryLeaderInitiatorTests { + + @ClassRule + public static final DynamoDbLocalRunning DYNAMO_DB_RUNNING = DynamoDbLocalRunning.isRunning(4567); + + private static AmazonDynamoDBAsync dynamoDB; + + @BeforeClass + public static void init() { + dynamoDB = DYNAMO_DB_RUNNING.getDynamoDB(); + + try { + dynamoDB.deleteTableAsync(DynamoDbLockRegistry.DEFAULT_TABLE_NAME); + + Waiter waiter = + dynamoDB.waiters() + .tableNotExists(); + + waiter.run(new WaiterParameters<>(new DescribeTableRequest(DynamoDbLockRegistry.DEFAULT_TABLE_NAME)) + .withPollingStrategy(new PollingStrategy(new MaxAttemptsRetryStrategy(25), + new FixedDelayStrategy(1)))); + } + catch (Exception e) { + + } + } + + @AfterClass + public static void destroy() { + dynamoDB.deleteTable(DynamoDbLockRegistry.DEFAULT_TABLE_NAME); + } + + @Test + public void testDistributedLeaderElection() throws Exception { + CountDownLatch granted = new CountDownLatch(1); + CountingPublisher countingPublisher = new CountingPublisher(granted); + List registries = new ArrayList<>(); + List initiators = new ArrayList<>(); + for (int i = 0; i < 2; i++) { + DynamoDbLockRegistry lockRepository = new DynamoDbLockRegistry(dynamoDB); + lockRepository.afterPropertiesSet(); + registries.add(lockRepository); + + LockRegistryLeaderInitiator initiator = + new LockRegistryLeaderInitiator( + lockRepository, + new DefaultCandidate("foo#" + i, "bar")); + initiator.setExecutorService( + Executors.newSingleThreadExecutor(new CustomizableThreadFactory("lock-leadership-" + i + "-"))); + initiator.setLeaderEventPublisher(countingPublisher); + initiators.add(initiator); + } + + for (LockRegistryLeaderInitiator initiator : initiators) { + initiator.start(); + } + + assertThat(granted.await(10, TimeUnit.SECONDS)).isTrue(); + + LockRegistryLeaderInitiator initiator1 = countingPublisher.initiator; + + LockRegistryLeaderInitiator initiator2 = null; + + for (LockRegistryLeaderInitiator initiator : initiators) { + if (initiator != initiator1) { + initiator2 = initiator; + break; + } + } + + assertThat(initiator2).isNotNull(); + + assertThat(initiator1.getContext().isLeader()).isTrue(); + assertThat(initiator2.getContext().isLeader()).isFalse(); + + final CountDownLatch granted1 = new CountDownLatch(1); + final CountDownLatch granted2 = new CountDownLatch(1); + CountDownLatch revoked1 = new CountDownLatch(1); + CountDownLatch revoked2 = new CountDownLatch(1); + CountDownLatch acquireLockFailed1 = new CountDownLatch(1); + CountDownLatch acquireLockFailed2 = new CountDownLatch(1); + + initiator1.setLeaderEventPublisher(new CountingPublisher(granted1, revoked1, acquireLockFailed1)); + + initiator2.setLeaderEventPublisher(new CountingPublisher(granted2, revoked2, acquireLockFailed2)); + + // It's hard to see round-robin election, so let's make the yielding initiator to sleep long before restarting + initiator1.setBusyWaitMillis(1000); + + initiator1.getContext().yield(); + + assertThat(revoked1.await(20, TimeUnit.SECONDS)).isTrue(); + assertThat(granted2.await(20, TimeUnit.SECONDS)).isTrue(); + + assertThat(initiator2.getContext().isLeader()).isTrue(); + assertThat(initiator1.getContext().isLeader()).isFalse(); + + initiator1.setBusyWaitMillis(LockRegistryLeaderInitiator.DEFAULT_BUSY_WAIT_TIME); + initiator2.setBusyWaitMillis(10000); + + initiator2.getContext().yield(); + + assertThat(revoked2.await(20, TimeUnit.SECONDS)).isTrue(); + assertThat(granted1.await(20, TimeUnit.SECONDS)).isTrue(); + + assertThat(initiator1.getContext().isLeader()).isTrue(); + assertThat(initiator2.getContext().isLeader()).isFalse(); + + initiator2.stop(); + + CountDownLatch revoked11 = new CountDownLatch(1); + initiator1.setLeaderEventPublisher(new CountingPublisher(new CountDownLatch(1), revoked11, + new CountDownLatch(1))); + + initiator1.getContext().yield(); + + assertThat(revoked11.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(initiator1.getContext().isLeader()).isFalse(); + + for (DynamoDbLockRegistry registry : registries) { + registry.destroy(); + } + } + + @Test + public void testLostConnection() throws Exception { + CountDownLatch granted = new CountDownLatch(1); + CountingPublisher countingPublisher = new CountingPublisher(granted); + + DynamoDbLockRegistry lockRepository = new DynamoDbLockRegistry(dynamoDB); + lockRepository.afterPropertiesSet(); + + LockRegistryLeaderInitiator initiator = new LockRegistryLeaderInitiator(lockRepository); + initiator.setLeaderEventPublisher(countingPublisher); + + initiator.start(); + + assertThat(granted.await(10, TimeUnit.SECONDS)).isTrue(); + + destroy(); + + assertThat(countingPublisher.revoked.await(10, TimeUnit.SECONDS)).isTrue(); + + granted = new CountDownLatch(1); + countingPublisher = new CountingPublisher(granted); + initiator.setLeaderEventPublisher(countingPublisher); + + init(); + + lockRepository.afterPropertiesSet(); + + assertThat(granted.await(10, TimeUnit.SECONDS)).isTrue(); + + initiator.stop(); + + lockRepository.destroy(); + } + + private static class CountingPublisher implements LeaderEventPublisher { + + private final CountDownLatch granted; + + private final CountDownLatch revoked; + + private final CountDownLatch acquireLockFailed; + + private volatile LockRegistryLeaderInitiator initiator; + + CountingPublisher(CountDownLatch granted, CountDownLatch revoked, CountDownLatch acquireLockFailed) { + this.granted = granted; + this.revoked = revoked; + this.acquireLockFailed = acquireLockFailed; + } + + CountingPublisher(CountDownLatch granted) { + this(granted, new CountDownLatch(1), new CountDownLatch(1)); + } + + @Override + public void publishOnRevoked(Object source, Context context, String role) { + this.revoked.countDown(); + } + + @Override + public void publishOnFailedToAcquire(Object source, Context context, String role) { + this.acquireLockFailed.countDown(); + } + + @Override + public void publishOnGranted(Object source, Context context, String role) { + this.initiator = (LockRegistryLeaderInitiator) source; + this.granted.countDown(); + } + + } + +} diff --git a/src/test/java/org/springframework/integration/aws/lock/DynamoDbLockRegistryTests.java b/src/test/java/org/springframework/integration/aws/lock/DynamoDbLockRegistryTests.java new file mode 100644 index 0000000..f0ff009 --- /dev/null +++ b/src/test/java/org/springframework/integration/aws/lock/DynamoDbLockRegistryTests.java @@ -0,0 +1,319 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.aws.lock; + +import static org.assertj.core.api.Java6Assertions.assertThat; + +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.locks.Lock; + +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.ClassRule; +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.task.AsyncTaskExecutor; +import org.springframework.core.task.SimpleAsyncTaskExecutor; +import org.springframework.integration.aws.DynamoDbLocalRunning; +import org.springframework.integration.test.util.TestUtils; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit4.SpringRunner; + +import com.amazonaws.services.dynamodbv2.AmazonDynamoDBAsync; +import com.amazonaws.services.dynamodbv2.model.DescribeTableRequest; +import com.amazonaws.waiters.FixedDelayStrategy; +import com.amazonaws.waiters.MaxAttemptsRetryStrategy; +import com.amazonaws.waiters.PollingStrategy; +import com.amazonaws.waiters.Waiter; +import com.amazonaws.waiters.WaiterParameters; + +/** + * @author Artem Bilan + * + * @since 2.0 + */ +@RunWith(SpringRunner.class) +@DirtiesContext +public class DynamoDbLockRegistryTests { + + @ClassRule + public static final DynamoDbLocalRunning DYNAMO_DB_RUNNING = DynamoDbLocalRunning.isRunning(4567); + + private final AsyncTaskExecutor taskExecutor = new SimpleAsyncTaskExecutor(); + + @Autowired + private DynamoDbLockRegistry dynamoDbLockRegistry; + + @BeforeClass + public static void setup() { + AmazonDynamoDBAsync dynamoDB = DYNAMO_DB_RUNNING.getDynamoDB(); + + try { + dynamoDB.deleteTableAsync(DynamoDbLockRegistry.DEFAULT_TABLE_NAME); + + Waiter waiter = + dynamoDB.waiters() + .tableNotExists(); + + waiter.run(new WaiterParameters<>(new DescribeTableRequest(DynamoDbLockRegistry.DEFAULT_TABLE_NAME)) + .withPollingStrategy(new PollingStrategy(new MaxAttemptsRetryStrategy(25), + new FixedDelayStrategy(1)))); + } + catch (Exception e) { + + } + } + + @Before + public void clear() { + this.dynamoDbLockRegistry.expireUnusedOlderThan(0); + } + + @Test + @SuppressWarnings("unchecked") + public void testLock() { + for (int i = 0; i < 10; i++) { + Lock lock = this.dynamoDbLockRegistry.obtain("foo"); + lock.lock(); + try { + assertThat(TestUtils.getPropertyValue(this.dynamoDbLockRegistry, "locks", Map.class)).hasSize(1); + } + finally { + lock.unlock(); + } + } + } + + @Test + @SuppressWarnings("unchecked") + public void testLockInterruptibly() throws Exception { + for (int i = 0; i < 10; i++) { + Lock lock = this.dynamoDbLockRegistry.obtain("foo"); + lock.lockInterruptibly(); + try { + assertThat(TestUtils.getPropertyValue(this.dynamoDbLockRegistry, "locks", Map.class)).hasSize(1); + } + finally { + lock.unlock(); + } + } + } + + @Test + public void testReentrantLock() { + for (int i = 0; i < 10; i++) { + Lock lock1 = this.dynamoDbLockRegistry.obtain("foo"); + lock1.lock(); + try { + Lock lock2 = this.dynamoDbLockRegistry.obtain("foo"); + assertThat(lock1).isSameAs(lock2); + lock2.lock(); + lock2.unlock(); + } + finally { + lock1.unlock(); + } + } + } + + @Test + public void testReentrantLockInterruptibly() throws Exception { + for (int i = 0; i < 10; i++) { + Lock lock1 = this.dynamoDbLockRegistry.obtain("foo"); + lock1.lockInterruptibly(); + try { + Lock lock2 = this.dynamoDbLockRegistry.obtain("foo"); + assertThat(lock1).isSameAs(lock2); + lock2.lockInterruptibly(); + lock2.unlock(); + } + finally { + lock1.unlock(); + } + } + } + + @Test + public void testTwoLocks() throws Exception { + for (int i = 0; i < 10; i++) { + Lock lock1 = this.dynamoDbLockRegistry.obtain("foo"); + lock1.lockInterruptibly(); + try { + Lock lock2 = this.dynamoDbLockRegistry.obtain("bar"); + assertThat(lock1).isNotSameAs(lock2); + lock2.lockInterruptibly(); + lock2.unlock(); + } + finally { + lock1.unlock(); + } + } + } + + @Test + public void testTwoThreadsSecondFailsToGetLock() throws Exception { + final Lock lock1 = this.dynamoDbLockRegistry.obtain("foo"); + lock1.lockInterruptibly(); + final AtomicBoolean locked = new AtomicBoolean(); + final CountDownLatch latch = new CountDownLatch(1); + Future result = this.taskExecutor.submit(() -> { + Lock lock2 = this.dynamoDbLockRegistry.obtain("foo"); + locked.set(lock2.tryLock(200, TimeUnit.MILLISECONDS)); + latch.countDown(); + try { + lock2.unlock(); + } + catch (Exception e) { + return e; + } + return null; + }); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(locked.get()).isFalse(); + lock1.unlock(); + Object ise = result.get(10, TimeUnit.SECONDS); + assertThat(ise).isInstanceOf(IllegalMonitorStateException.class); + assertThat(((Exception) ise).getMessage()).contains("You do not own"); + } + + @Test + public void testTwoThreads() throws Exception { + final Lock lock1 = this.dynamoDbLockRegistry.obtain("foo"); + final AtomicBoolean locked = new AtomicBoolean(); + final CountDownLatch latch1 = new CountDownLatch(1); + final CountDownLatch latch2 = new CountDownLatch(1); + final CountDownLatch latch3 = new CountDownLatch(1); + lock1.lockInterruptibly(); + this.taskExecutor.execute(() -> { + Lock lock2 = this.dynamoDbLockRegistry.obtain("foo"); + try { + latch1.countDown(); + lock2.lockInterruptibly(); + latch2.await(10, TimeUnit.SECONDS); + locked.set(true); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + finally { + lock2.unlock(); + latch3.countDown(); + } + }); + + assertThat(latch1.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(locked.get()).isFalse(); + + lock1.unlock(); + latch2.countDown(); + + assertThat(latch3.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(locked.get()).isTrue(); + } + + @Test + public void testTwoThreadsDifferentRegistries() throws Exception { + final DynamoDbLockRegistry registry1 = new DynamoDbLockRegistry(DYNAMO_DB_RUNNING.getDynamoDB()); + registry1.afterPropertiesSet(); + final DynamoDbLockRegistry registry2 = new DynamoDbLockRegistry(DYNAMO_DB_RUNNING.getDynamoDB()); + registry2.afterPropertiesSet(); + + final Lock lock1 = registry1.obtain("foo"); + final AtomicBoolean locked = new AtomicBoolean(); + final CountDownLatch latch1 = new CountDownLatch(1); + final CountDownLatch latch2 = new CountDownLatch(1); + final CountDownLatch latch3 = new CountDownLatch(1); + lock1.lockInterruptibly(); + this.taskExecutor.execute(() -> { + Lock lock2 = registry2.obtain("foo"); + try { + latch1.countDown(); + lock2.lockInterruptibly(); + latch2.await(10, TimeUnit.SECONDS); + locked.set(true); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + finally { + lock2.unlock(); + latch3.countDown(); + } + }); + assertThat(latch1.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(locked.get()).isFalse(); + + lock1.unlock(); + latch2.countDown(); + + assertThat(latch3.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(locked.get()).isTrue(); + + registry1.destroy(); + registry2.destroy(); + } + + @Test + public void testTwoThreadsWrongOneUnlocks() throws Exception { + final Lock lock = this.dynamoDbLockRegistry.obtain("foo"); + lock.lockInterruptibly(); + final AtomicBoolean locked = new AtomicBoolean(); + final CountDownLatch latch = new CountDownLatch(1); + Future result = + this.taskExecutor.submit(() -> { + try { + lock.unlock(); + } + catch (Exception e) { + latch.countDown(); + return e; + } + return null; + }); + + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(locked.get()).isFalse(); + + lock.unlock(); + Object imse = result.get(10, TimeUnit.SECONDS); + assertThat(imse).isInstanceOf(IllegalMonitorStateException.class); + assertThat(((Exception) imse).getMessage()).contains("You do not own"); + } + + + @Configuration + public static class ContextConfiguration { + + @Bean + public DynamoDbLockRegistry dynamoDbLockRegistry() { + DynamoDbLockRegistry dynamoDbLockRegistry = new DynamoDbLockRegistry(DYNAMO_DB_RUNNING.getDynamoDB()); + dynamoDbLockRegistry.setHeartbeatPeriod(1); + dynamoDbLockRegistry.setRefreshPeriod(10); + return dynamoDbLockRegistry; + } + + } + +}