From df4701957ba0c2c4dec9e69a2b3aa5567127fe13 Mon Sep 17 00:00:00 2001 From: Michael Freno Date: Thu, 13 Aug 2026 19:34:30 -0400 Subject: [PATCH] refactor(feed): port refresh batch to Effect for concurrency and timeout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace the hand-rolled worker pool (mapWithConcurrency) with an Effect program (src/effects/feed-refresh.ts): Effect.forEach bounds in-flight fetches, Effect.timeout bounds each feed via the Clock service, and failures fold to null so a bad feed never fails the batch. Per-feed apply-as-it-lands is preserved — the apply callback runs inside each feed's own fiber, so there is no Promise.all barrier. Store boundary unchanged: refreshAllFeeds runs the program through Effect.runPromise, keeping runAutoDownload + flushPendingSave after the batch and isLoadingFeeds around it. Adds TestClock-driven tests (tests/feed-refresh-effect.test.ts) that pin concurrency, per-feed apply, timeout, and failure isolation without real 20s waits. Pins effect@^3 (V4 is in beta). --- bun.lockb | Bin 78110 -> 79571 bytes package.json | 1 + src/effects/feed-refresh.ts | 85 +++++++++++ src/stores/feed.ts | 89 +++++------- tests/feed-refresh-effect.test.ts | 230 ++++++++++++++++++++++++++++++ 5 files changed, 349 insertions(+), 56 deletions(-) create mode 100644 src/effects/feed-refresh.ts create mode 100644 tests/feed-refresh-effect.test.ts diff --git a/bun.lockb b/bun.lockb index bad57d5f40e5866d9534bd5441ea2689e98e782d..83b76e0aea52d08471b011a5512f4f53d0176f5c 100755 GIT binary patch delta 14516 zcmeHOd3aPsw!in14(ULa1`;l$fh0g6f$SRuLW03Yxd91_LG~u;>q z?yXMs>Pv3NuDdOX2(`Vr>hA+=Ia5!*o9$D$&Ogb|XY+HvCT;jSqf3;ut=}5|v;B-Q=i1M^26-t6UIfqr^`TT+={3L0d_JU08IpC{bryt=oQrG zJ}d1-MfuKh!RaWT=qMAaggz)}hk_5$S}^FVpjOc9Zi3(g`W>ud^5I9Q8^UL!c&c3LGzFm=^ z9DyOt1)(*l;ox5paqCbB${fz|MH3yuRm6=WaRHP=dlZyA%(FYkPbv|F4v0H92z^Hp z2DY>?H@u=$5OVTMvN_`icXcD z7{}cidNa^H2RsHoa<>+Y3hVuh=y9Mgf^vZO>-;92Ujd4ys$5mbu;b- zWEBd6v)rER2rqC}pq%|&4X4qUYo1P>pxp3iP#%B`P&ObF`aC7mLD{ilds(3)TM!EJ zixCAuxZYlqM|I$KF|wC*!IbgPLOKAR=k6#__RL;V89vVLbi4{4gX2oY{9}WP^NZPD zp4|IZyJ^lH39qQY94WiUp;PT!dU0qAQFmE)8c>cj~U+hnz!F2cBb|U(91PLGXyw z46TBkJN^Vb(|17Gp~=NXG2GL(C@pkO`((C5C<5=tW3U#A&7oMQS6r%hRDuS%!*rc? z1LaB35tQ|PLD`V&(B}?MfpYm)P%dAqw=))$Jsw}2o132_2>6&(W-k@`LGD?Ff*tXi zWqL@2?ED<42|_YP)Cc7lA$8EtNZM>IO%w#kTmwMShU@EYf`Aaa{sfBAarI3Sgsz}L zpcoZbDdapWP2hRoS_B>wR^3oTc#PlZsd)k=wWA)%+MF;#`%W@VagPWb6L@}E^?{GK zuRZYO{=^lpJyYH)=JH*e_q!cyK3(@g;Gy@oY;aF+>;-3Bjo5r7e_GQ#(Yy+qdT&TY z3Tl^PISUT$KhSAO8Eq5I!SD^{x!hUDJQ5s?bh$AnEt4q8+a?7yqGv&dH=@tIZKjve zM*@B6ohF@VM8kb-lDixI5v0tGuKL(av#=#*lUM6BDcqf=_}WbK-37r$*L~8YBktsB zv6-hd7K9GK|c?jEhXdO4ABY`Q>P!9@gW|LNWP#MTo584IN(@auxn|V1l zCN|haYb+_|L*Np?HDXR`=1HfU+e|jB)Uj*>rTV#1QHwsNeu(yPmY7~aE|twO-9av$ zzDZ3pJ&sw@hh^ra$Z_LtG@(t3bf77<_P3e(VmZZ7n|5hZr5C;KZpzG3$U4Mu}*kg-KlqIinOO0?Fz7&BrKLeYEVj$dzy2WeaH=mBF0Ef&1iyGqSUqp zownL6&tsKwE0{=X@~1lQM9J(&1#N87ct6?<@`fMXZDTXv!~%_jszev9DduQwt&ed< z>fJI$a{ANJw$Sva*6nQO^C;6s5Fs}QVt4Na%|`6GIsYMUJvi=OD?1O4>q!(HkRtiX z^m(Ao+*cL^Y%eY|Rkuo!{vgxvAe;F|2zx;2Nw?$_a~K{_+)87%#WW6FFDh=6CcWH> zrUcu}e}TZJHp_-5x5mr>C$bxsYH(cDP4^$pVFm4kw4Yjo8*GHVBZ-o+JOZs0*1;xL zlGw4Y=?kp?0W3DPYQs_ITp@CrY$tLRae6zNP5C575dzQG73eHof2{oI~Dhajl~Pi7WO?7#x)tq%0ZFq`Qd ze`Xl5pw2I@nCUV$3pfMO5B#)Fjn9|G;dcLd8{2hd&Rbh2- zij>ib!oqFlZFoTTQQP8ycI!;L!fn!Hok@zoUV`nIgVcy8f#ofo!;YgRUION8kn%(k z3Eo+PLbVlc%wcRJ*eORq`h?mJPrI$)ctlOSTFm#r@sxDq`sQ?O^t=iM?$JCK9JlSR zatFZiG}m@d?=I9j+Gg(G#Rx8%HqB>F4aH~3**GWq(3;r3W(B(euV$_C^se-Ij7@xt zqT~9SujoaMc~O{Jha2PYbU0m&wV4(`ki;!XpM=wtIGgD%1X{cwi_k~JG#j~eHS$N1 zW1H1I&eA+m8!7BKYBF;?ObVo`apsfA@m3`ACL?u-qFo6#b2Uazott#TJH`9~I1Wh@ zHC(@dLl3BMD;BFQ8dFZ?5K7ZDa095gTblHFG!=BSnf&3gI%7=3ALLd)$X!CNpBkdB zbStT^sXSKq!L$cC)jpFu>`-+~1Ci6bUWuI6>RIHp(D=vG=RIr|1y9LD*rR(Zc&hF` z3#m3O498h-D8rwVvLvnWvuV7f1)G zPlA=`v`pP=J+Urw)SS6H&R=TjV&o>OGG~$?OrhcyX{PJQm2jzL1U4^T6~>+%ptije zQtgq}9388VJIUFTPN&+;dmv!bw5`!SnF@N_Onop!NmSfB%{&7+UYxDe1@R#`-b#$t zeW|)>Y!!11exdLLc)8$tsiuaMA76Mps);G$^5wBgwP1=VG6(=AgF3$%9cK< zu_|SQ1~5l7WNB+tT5>t-M9f0M(0&K{AfZNqRc#44I(0Qs-jk01>@JzjR`TvEoK^DDUbDjEu zvRnqmpU_I@t)M3EKL{DNv=b4Qplo?0C?}>|9<9scLAgl(D()f5@9kmxglh~XU)!` zldin}+5n4r;@hi!J~Z#_%(crt+`YEbt1;`7c5FVTK9TPE+ymvOkdfTxytJ)H&+{JP z`>M8$*t&k?2fKgReB5_$>g93Eo4=p9Y~zxy8-H%{($?5k#k+jAzF)X@K+`w+9ZLEr zcEzrpvt}RJQq<@xO3skei-$h{;oVj*`-iT7vvy6Bfg>kotxCT0c*njkrcvdcX+!s4 zyZza*Sgf z4N*itvJc55`Driu7+g!T49%qd;3|eHB2H}JCTDt4$EOuBfF?YhNu7sz(Q$BXC@3?N zPJ*k+RK#|41YGSiUKBq}5re3DSSEEF?nRfuwWsK3GU*Dq<Qv%^1diDVfI|G-s@RmASJ7u@7=@XxM@J!pa* z{$;^GaLE)j4*r3w8K;OTbOc;&HvG#{#NJe$1^*oI4_q2WXTv{m%d-{nak>a@We)sv zC}Ll#bHKk`_y_Ju>YW4sz^%(sM1^jF+mHwUauu;Zt;vOd`S34K5uc)rJor}t|G;IC zm=FKJIr0^8Fl_>tR|x+K6mclo3o^y0X(v82$x@gp4x@a0K0|x)Ihu#e5RO9bMfh*i}=i;?xmSxF4dL7 zn=*J)s)+g2djh-xw{C(W7Sb(n8=UZ_Oc9G|O&Po?hc`||98VcecvAsyz)c{r9NvI) zlq;f>Hi65V2yZGBv4ZRs7_CVdEpU^_G7+N%u41BsZ#MUWn>-n#HAxYt(u7GEEqq0O z9NcsYnvBr`S2I}=tLO;0+9?>VN=2-u>Pn2(RE!q5nG`(*qXlmH6h*9|i{MsH!)Q%a z#5q(q75+_!f8gd)?`iN4+`4IsIG=8T+b{$EO;^N)v}QW|tAc+s6mcly4mor7XE>Ije5_4f8f^5QN%as7Pt*_;a{yH{)yJq!oPX&Z>}P~O&N3H-+cH7 zZWD>~;2${0JVo3>o51BQfPeE9aU0p^!@q^_58QULEP#LDDi$c>PTC7@@*?=RP!Zpy z2@B!hV)zH{JqlU`|G?EOQp6AF$f8Vf4}~ty6!%g!K0l;W`21)~^peFQZCo;te7)TK zjf((TMfnl@pFchwWH+YMbrZde|Mf?kae$z@C0)1zEm-PdGU0&X1<6`fl1d(?UN}*h zA(=`8L*1lBK6H>?HO;f&DDprpaJh$RaSJrHrEW>2xFvPDRZ^=@Tj^ms7^F4$1xvOJ zd#={a%lJp%hd45LqmdL^66~=wt}`kQ5w~o7X-Nn8u+MWPAHb1**-!oUC8>vvd-t`( z;++6C{|Db&89cig#|22c0M7F9fDf^p__)t9K3hxyIPC_wUgyu46NJD68ujd15gXGh zFfP!UPxtqQR`Jn~Pi!@y^MJX)0$@I{2v`Wr2Ic^@fD5Pwo&{zCvw+0_e<#PMKtz%M zXUhL4<^PSUhdF-6@`(~dqn-+vffBHs4}X6^hQBbWqir9Lsd@{!w}FkoCLjlqr3Duz_Gs#8V(u(+(7;s@IBxQ{#(#}fC+hbpfTV9@c8hlxBz$pc|O+l1^NRK zKseA9z!$N?(3|MQo1y-^_z=mR_kR6;fxm;~?uPXLO5VxR;V3)q1pz#-r;a0TEm zF`57akskm=0ntDV5DUZse7NloophjI6+WgzPzH<#vVfz&d4RvW;qP=qys;O_`{_VB43)zq2Z2$W|yQ=Va5muK1gz*d0k^#XWSdIL=WULjsUQ=KM)vfLNo z;3NQT0r*y>4h;_zM=25r210-i07s9b*bcA&T*(Tw0Gb2M06wy{26%qSfFIBj@CR70 z4Zv*%0f7J;!P9|dQ9KoxVFS2xqvQ6-cLcbRK!D3xuM5x_=mdlTp*mQ`^DGcxJ)=!_ zf~SH}kLBD?GSCC)ivHD!$@7CP>kjnP^Q`ye?Jb6=q$AI2IDv$>B0eHec0rUf& z1QZ}07zzvmGJpZVK;S80FfarN=dsufYyt`aUaq`s^MG-{SYQk=8W=^+{YYXHZT#W0A}%w zc@~*Tz+}J$@LD$3S0(bidGNr_0IGF56SM|kYIMdOG<1Lsm=7#yptGn!ehKHfVZ-vJ zdLf6fp+|=2h92{O0GO@-*nq_V4-F3s>%9!T1UwHs2e9sHU={Ep@B+YkRctjEGV@#7 zu&iO{++jl_xgN{7^WPe7w9Dn(9#fY633!tO%bpBN%1FRbX z!~u~2cM<}00Ahg;05;MBbOAa8dw~yukAQCGqxeM1{)tY1xj{B=fEt&JBylu1BarLh ztb%Exx~aR2n?_t15fd8`{oul=aTf_C(GgJ*u-3RbYTQafNkRn5@Vt}MZQQl>PTSK< zT&s6TVq!#0H!e9NQQ)B%t9suNOCsv+GuM_5@C%m2c(fSYJ|9FeWYmCjVZX0Ki9zU)<49}^K9fynnUk?-vm?%g61ge()CI}&3x zu1((DH@)UvbR8aHDW=+}q;>*(|F|4uXXX}8NCpFMo|-gjzM1S^qW zwJ#lwkq$Pd!(Yh>#>@Elm4|$0)TXcS6;GleGoozVv%S4ygu7yC^h{WlNURdjFJ5*eesXc z-?h(!g^3Z-JdA1;S?b%2hJUTDpG6=;n$v;nGFH>l8?x27{OsTS&*`0OpV(VpcSZ~H zJ?XLT=*^R|bhQP2axzBh;7@l!{901PsTc`2jmMqJ zNic2_x05c7&bV=dhZGB+m$q@&*!2DKgp!qSHARWOk<6Fr-l?z-$Kg?H=sjY!MWY48 z|8y%#KP^jx=ePetPhgZ|fC|Yt6s(uMA9iBg~_| z^MzFFA!4 z&dFBe60i4XFJwmNZF{O-;eIeJxGGD&?P(*w(!~F zYe$azmTjq57}cJlV1G$_I(JS^Fs>_4cAN8>$4@y?^?Jqy=cSGy=ZT{sSm%ZpuePVI z=VhyLBie3_^~!S(`?Ox+KzkZ{UY0Jkr#&$FdVAXV=LcL=2iku=M$*Q?`i}#mo7u39 z2epj5*}os7M-0wC=|gkmVLQonv9)fURuE$~?v&?$wl!q)%3pt}ceJ5(a)NOOeb?pj zz1%OHHen6(lLWKMxTW52O^kog=G1bO@Knam+b@KQE~O_J*WA4`zFYA>FTCx7LY(#t zHLlc8s(yTVC@tfvdZUkZq8pcFX+GNy^zFHF`H9*V?+p?qiz_3)AR<+z8n zY2xLFBPr9q-%9k!L0$j+KvkdUj~E+c{2v;y80zWfn!tGY!!NhB3qiIIw;X!VtJXYU9TQ>uv=7W!1cOKN?k| z@ur36wDC&CCabS$F~(~xS2iV4+aFTGjT!b?VQqqH|H&2s&H z-}vc;kn{e&>q{n_g>HN7jmqCea%3s1h`$Dp+K(W0(3?Ag*V)AAYMsKcy9Aj_@*cz%2%}e858g SK|VZz?*Fa*maBh@>-oQVnl;)0 delta 13725 zcmeHOd0bW1_CM#yMJ_U35k27*0TU4g;R+XVLPc-t!CAbH2*PEONf{I;u9|bAxU$qn zU-h2l%Odk>*~EKRYFJuoq*jhkn@p{KmG$iRU1tz9U%&VI{oDTdJ>R|dxYpXk+2^wT z;46H2J+{k9Wzjtj-HTYfmV<~{3`UmglQdNsWJ@P5l)pBlgWgPK@T;P6JR zv!gU%U+vu3+WoQB*OAy-5Zrsgt-vEa1R)H39UKC_1m^rxU?1=(PeEu4ehKUic9msi z7G`F<20`|P90P6x4hM7jf0+b927iG1W^h$zVPSq*xlraRs&vg3s)ZiN=!lGGksA)K z1zW+5Ed;?I{2qD<0)G$L0xpM^+|WhH0pNl%r~tD9(A!;*4}#l+cY@o2e?)oaKcPG; zU z{ID)B0pnk_y8ur%ygYBV%Y{Zt3p0zeW(q=Cd1kIFx}Yor`E2JZ%rl-c6*{Mbx#3}8 znC$KgW(5YJJfiMIFdJ5sIeVrnTM!EJi{J%8XvFm4^!A;&U5xA@Jz-u6N+Ilq%yV}r zm@UgJu8N+KS>~#R3}fBhF!5QTqWmJ(S4ilrdFVwjPpRk#%>x_2?5Xo$?yz)rNtUau zOmJ0V?8^jUMnz7Ji;6s2RUhi2wXg&wd9v0-YCS)V7P)4Zu3CCo`RwwNnLUMPAhW4a zdVWWpEnx2P1!%}~u@ubZD$4V7czjATXO^(W*+}O}URs=607a?=r(Uq2)RhaD3(0zc z8AYYhtm>1x>7OH=U9tc&`#ish$7q)DO|+(H0n)kSx)_bOg4v+CMTIu*>G2+#>&i0c zvL1pHauAO}4GIQw!+OPqdPkGcAa|Ikb0;uQf&eg={|V(-kw!3gybsLzkAXSA2F&#p zFk4(wl$)EMBM5lRnVnfGBp}_p8W~T-YntgU5wi1hP)raaF^hOQV?@+Je@oP6>l8ck zke&pFF7Cf43Ibg0-UN;U+mZyKJJ<(|QFTv2I?u{0y|jI+3Nr6!yrD+$7(dWQvjm0I zjw1SMbHZ@#*1ojGGiK){&(xMCwBh!k#eT_BwMbk19g>ekXF)?GiVAR;)<}Y2r|5tb zsa~R&102#-i6o0dn%07bTO6hm>{QuQ)ILR$O>_bYl_o)O&yQ6u~h0e zD8&>D7Y*h#^FqY9aSu8Yk}U0RMf2J@%r=ZId%uP1lx!dJmL1Xs9~uiv^rbDJ#lCb_ zc9dhxrM@MEV%r)F9SZO*@F#0M#|#{xm$;Vd;yp=V63L)gV%R>qIHfpH2ij z%sY^o%9*Mv*ZrxWy~C>)=6gRR(i=W99c~|Ku0{guCb2Uub&&EQ>Fo)TYJ{h#h$b^Dp}fYq2LY<(^p6wt2$;9X7~h-S)M~|G74fy)X)%T6c6v5454l9UbPASSy(L&~0Eyve}B2G>9wGRykQp z4wJ%C0!|yMY&%z-C;U_qEBIH>%>k_sZE$|{RA4e+B5s_e+ z#Rt19uXAG-n`oe4qS=Aep=!rGs~&^IgK6URV*U&g&qfcfZ%)KJlGmZY+L+5Baoe7% zv>OsnbZrm)u`A89In2GH48NjjQzc8PD_%s*p~kKwrux)u@20J3t%1C56cz6;CKEOgwVE`rI}J~8h;wLG{5?`#cRG>aFkL{d=JRO0M70r_g;=WE(yNHE z-s)~=xdDl%H1-fRnAs;LJKmY9g64l9#v78hZ~OP4;9d@MA>OgY&1y1S~d;^H^1aJ59NeM$)0A6lqs19ZYhV#5hC0Ava^IZpIEFHe7X6JM4Cv zMbi;e?J#XeOmpH@#ME+PAByX9k7;R7qt&+&Q~f82@f6k1VX;GPye8ecwSu>5ZV^(o zsks+1yd^RGJjCBaLfaT(o*G>fv}x<5?rf7F@jfB&He%ThNv)?hZ06}_x`@=V>ib|R zJU(09alb<>M~!_8MV70%^DqLHYFc`dAk3pfK`Ew2#ELoB;=r8cm0|46?bWuQM5^|3 zyMdVc;*-+*Q13ww^LA`@+7`yE=4u~0ILKk@0rh&*p+PC;e8hOj?@|}VvygZvFo{+04^T{aA?k)F9UeZ30Ru~jtTrSR6}z{ zdbNPzR~y7AsNuh1Rsh~tLv!{(3OuU@W_}C9T>=|`^{WQv7CcRc=n2fM>2O_U=89Mt zYG}?9br{mFn!6tZu)X5}meT+Z%v{^a0s=G3lUP7#&O8Rw0j{3`aL7=f8gsoFy3EYw zvUtg{x_JOAFcaX=oVj2jz+41yXwIBo3^11f5}6zws}T{(0QTY>fV-~(INXW3_B?=j zKERC9a2gf26)d`g#pi&>E~04sPN;C8Do;K6~JGcK`!z|8U&EFduR9Qm3B1ZMWkcL1mV z2(ThQ0o?G<00(CFoF^yh%;mfwgR2E^J>x%NR>)VYDERAa0dsmAF#ZWax-5g);tpU| zv=f*+?yPecum_KSH$*rvbAf0zLlEM?oMq$3ZOlnM^>k+LAPLNh^ws4-U~X?P*bAJh z%Ok<~CyZj|@gI!{yKp?1D@*|UgI&6u3&uYoPv?A{3ph~^tN?=*_-}UiYSne9#BBuu z9GWxd|A*WA@7dk8`OmKbRv-c3(45&bNouOb?4druuebN#v%AB*Kel(>yS&}=DEzU# z*TV}u@%I6LZ0~<;@4tI{?~4iaySMk9ca1(O(ui>i^&0C-PmNQ=mUMNTlP*C@8?T6F z+B)7zo5uN)DNPYu(U>$R4IS@GyCC_JIKfFjKypn`M1QJ-v@Ok-f}DzIp-iWfCQk6B zHz5U*WulY3oW4{sQ4xcv9?~92ohK<`JDN4gNqG}}=@=w@T}*b8Jjs`qO;*JAbOh3V zNC{IEu>&oh;-tBgedz+Ejubc5Ns&{0Y2#Ey45KrU8X*mtrih)WcAAq`PW7elAa$mJ z=}zi34fdrgVi&p!=@O*0>5ABuwoZqA>98+D5xdct4A?gv_Cbm!F%$Mda%C!F52}N- zEd%z=P{cUOoB{hXVIQQPWXXbkGhknqA|_Bhq&<*2XDebN&B}&-S+EaM5{0>7UpDM> zDfsX_0%<>_gd9ceM~ib{p9}UuN~XA6*p~zQausnPoq^N{X-J+TrciAj?8}9HkcQB} zeAt%<`|=g>9=ZzY5~Q>OMI1(33t(S9?3<~G3XPcw`wCzmq!A<*!ahi@LPZ=!b&$5r zgndPdIEFHdU|%8ZgEWpT#jvjk_7y8)8r2s&#R+s*iBojaEIcRD0X!#BSgBK-OjUSJ zp(A)srLMD_;xt-}XF7d^=X8pj?G!U;6`q-N2G1FkROS@3s20y``V3DO4J>zxIkXwi zT)J8gYsz3vg(BwD)(Tit4r?kEaVCwagf$hg22v4;b6^c5*BnJGp*l$0Dq+oBMVv*M zb79RKSOckyELE^(F083i#0sj1vc?$m`G7qCwh0%gkMPc(XTJtbk^A&ME9f7nT zQo;g7Tu6%-V6^6Av>;Ve+(L}j0*ux|MO;j0AT>f7vPcn^Qtcv))|3meYiP`3*yo0Qkk*p81olC4Em6dER0nAreoGEo zs))6exfJ#-fqjsOEce2`rLgZ_MchdBkoG|8yi5@vq*=>g-@UL8(q;-<4*QnDzU7Me z5FLTEA5y{!MSO%7uYi5aVIQQ&C~hU}TLJr4D&kf;1E~?xkX4HK1l6vBeJf!fq^D@$ zYS_05_N`XLztB}kmmsCBQN*Wd>l)a%8us0%h}&t*eXws0?1S_yiECjWB-dI+e2(fM zZMzTl)hJ>eW!AvHwXhG;PO_|neKoLeog(g{dPsX9bzZNCyJ^;X*tZV$L3)Y8YGL1c z*jKBFd+7+I{g4tiDB`QMcmwRKg?*4--x){iL>leu(bo9OaaLh|4F9Q!4}+sTDQaV- zukpuWE6&=jsFu2M1v<6S%X9@PzDRjOO@8 z@MkI4<2Q6Az+n%-^#a~xNf3T`TVwCO6+LbJL*qls^5?lqawxgMLfsFvsphYX#VE2I zSO%;FRsgGkRlpKpDR3{a09Xht0{GG6YrpZYO>KZK3I19afP^y(QXo&ih;vVkmMDli3@2BZU%flPoumL>sCU?LC? zBmn$XmI&B^EVP^rxPTlW7YG9QLuwzGJ+vP<032k?*yooabOlENpCf(|I0yJZ{uKNM zAR;aSEdURI$A!Pza)7~z4+2tvVL%kn6^I1zgCYMyCyWHf1N~4w85jWc=WoIaB$fl1 z;_AQ_f%AZTpa7T(Oal%9?*I+Jd4RvNTLP(w4+pvf(LfB)1BeAMLxo`|GZb(DeC#g< zW&+cJcX?z_0{q(o|K7l#>SiDgasDk3cnn2Y?5IebEgF0m6Y$fPKWiY7h7VTt@~1cvko$;s>+?+5$X< z+5i?H5NHjcAYW)u8|?tL0;~jg${n#Xh63EV(Q!w_!vJoB>#%}ct_yHC(23Vu1R|Yv za3W79F3bguMne(j5jQGwI(L)=Bmyk+^y68;iX{NO^f>2@1%?4Ubq50csofvw$Frg@ zkPHj}cwP?%9KaCZ9$+Xi1{ev90C);RqxasghDyfCwYbYLnl1(*!* zteFTnf$9l(qyggr-a2{y^U9t9WC9t$bf6S?23QOH1y}>D0+s`h0DlIa0+sXe z?U4TpJO}Vn<7MRsya3bzJ^-h+Zx-K;_)EZxz)n3q5*z~X=EJ-P*v;erH$-*;w-@9J zMkS++(HV1hzy@%GaR8T%0%Cz~0Cy4&gaSQ)y#Ooe3v>bQ23`%q<2B%QK>d9iwC@_V zet#PseSeA7xbyeok0B@4yC2;niHR|`UNPw6Es@^%z-CpisA2BLK6mnq4I_iXB{3l; zHip%@E|TBT7o~nJXx~wrv{$47-^vtlEW>+a3snbPx1-Z1+ju6SwH=6lr$zU+pdH6# z=|Bs*+$c+pEvWPxE#qPf`WYGDx1a#MmX`;`9+$1gExZ9it)D)dbpM|vG2Ir^GbRC6 zd(iX$)5fz`j9nP%L5q&ttj0~ZEBh9fKKxqek5Cs8<&^1t@b~D+_My&#X~W8MYHKQv}rcKh5gu zuEZ$)^(5x5MY~e=ed@ZeG5l_ z)wr+PHt?URkt>JRH&xCLq~}k$t;Y4xCkiSns%}i_fC9FdI6J2MJAst?u`C@6q@0g! zl7A3w0SRsCT@WsZN}uG|jmw`MrL*Z{E??#`#lq*MZCnpEUEF9d-n7F9Iq*8WwXz*O z`bpH?Z?sc)rJMIko2|Auw1ARt%Jk(YvJ`w5d46h>B7>>-r#V*R*5=0-`nQ=@dL>^H zd!jIJZG(d8l}}|UBbbhVnrby}jczO(6>v2@qOFG*i#42x&3Q>M^*=4!jeDlg)addPcSycU!;dJ1kJ{UvaUiK({bjkG9o&ifG~R;A4F*6>O}Q#JCtcd;8{9%{c`l z^Y=WRi2cCGQNrkxGqU6iqrwZalpRLJU&>Nx7(EVYVHoZBT9#_UDC4YbH*QNFn>Kmn z&h0BPM~pU%o1_tMHV$99;2*P_a`vHxv$EB=c<2bgsvW zr@g+;1_~i>F{XI{gFkun%y9Cq9xLG`WtF3L=3j@lL!*dyX zZ$fACxtMA?F)fP+H(b7|t869bYyvbOlY-#gO+` zHp$t8`hS%wJ>G-v`|9Q-R;RR0pB}f6?jk+!_6oXscd+oLX}Yp^v|hE~W=VZr>1}G( z(ZMTu(jDr`(i8Et4!7wIi97WTz0k4UJpb zV?N4!Be5ajFHOa?*QN3DY<4KTCBtf5(|#pMeEWUZM?V@3^c1QR>GCz%ZrtzoGoKz8 za(VAGtpLBQjT_)kT@L%~q18`)ZB&ip!xmoE#<7Z(R*z~n<7CUVuk@nVzJURFf#GH0 z{&oL8Np$L4oCT7|=UdtOo1Nj1ZFJTPN8p?1!UXL=(Z-Wi&`&jqF+H0K*o9kq(GTo~ z(=;<^>UY07zVpX`dLEMPdvUPGmy!FPvxP&V>UY0nc)>0h|AA)EsI*Zd*3xT6j!o}B zlv^?WyUrH{^Ou`@+aELB8>=0&G}ZMpzj0$;$a!twSL`!9BKdn diff --git a/package.json b/package.json index 1f7172e..e958f0a 100644 --- a/package.json +++ b/package.json @@ -24,6 +24,7 @@ "@opentui/core": "^0.1.77", "@opentui/solid": "^0.1.77", "date-fns": "^4.1.0", + "effect": "^3", "solid-js": "^1.9.9" } } diff --git a/src/effects/feed-refresh.ts b/src/effects/feed-refresh.ts new file mode 100644 index 0000000..562c564 --- /dev/null +++ b/src/effects/feed-refresh.ts @@ -0,0 +1,85 @@ +/** + * Feed-refresh batch as an Effect program. + * + * Replaces the hand-rolled worker pool (mapWithConcurrency) + per-feed + * fetch/apply plumbing in stores/feed.ts with Effect's structured + * concurrency: + * - `Effect.forEach(..., { concurrency })` bounds in-flight fetches to + * `concurrency` (starts exactly that many fibers; each completion pulls + * the next feed — identical semantics to the old shared-counter pool). + * - `Effect.timeout` bounds each feed's fetch to `timeoutMs`. It runs + * through the `Clock` service, so under `TestContext` the TestClock + * drives it deterministically (no real 20s wait in tests). + * - Failures are folded to a null result: a failed or timed-out feed is + * left untouched instead of failing the batch. + * - The apply callback runs inside each feed's own fiber, so a feed's + * refreshed episodes land AS ITS OWN FETCH COMPLETES — the + * per-feed-apply-as-it-lands contract, no Promise.all barrier. + * + * The store boundary (stores/feed.ts) supplies the real fetch and apply + * closures and runs the program with Effect.runPromise. + */ + +import { Duration, Effect } from "effect" +import type { Episode } from "../types/episode" +import type { Feed } from "../types/feed" + +/** Result of fetching one feed's RSS. `episodes: null` means the fetch + * failed or timed out — callers must leave that feed untouched. */ +export interface RefreshFetchResult { + episodes: Episode[] | null + coverUrl: string | undefined +} + +/** Result guaranteed to have parsed episodes (the apply path only). */ +export interface RefreshSuccess { + episodes: Episode[] + coverUrl: string | undefined +} + +export interface RefreshBatchOptions { + /** Max simultaneous in-flight fetches. */ + concurrency: number + /** Per-feed fetch timeout in milliseconds. */ + timeoutMs: number +} + +/** Fold any failure (network error, timeout, rejection) to a null result so + * one bad feed can never fail the batch. */ +const failedResult: RefreshFetchResult = { episodes: null, coverUrl: undefined } + +/** Fetch one feed with a timeout, applying its result as its own fetch + * lands. A failed or timed-out fetch yields null — the feed is untouched. */ +const refreshOne = ( + feed: Feed, + fetchOne: (feed: Feed) => Promise, + applyOne: (feed: Feed, result: RefreshSuccess) => void, + timeoutMs: number, +): Effect.Effect => + Effect.tryPromise(() => fetchOne(feed)).pipe( + Effect.timeout(Duration.millis(timeoutMs)), + Effect.catchAll(() => Effect.succeed(failedResult)), + Effect.flatMap((result) => { + if (result.episodes === null) return Effect.void + // Capture the narrowed array before the closure — TS drops the + // `episodes !== null` narrowing inside Effect.sync's callback. + const episodes = result.episodes + return Effect.sync(() => applyOne(feed, { episodes, coverUrl: result.coverUrl })) + }), + ) + +/** Refresh every feed with bounded concurrency. Each feed's refreshed + * episodes are applied as its own fetch lands (no barrier); a failed or + * timed-out feed is left untouched. The program never fails — failures + * are folded to per-feed no-ops. */ +export const refreshFeedsBatch = ( + feeds: readonly Feed[], + fetchOne: (feed: Feed) => Promise, + applyOne: (feed: Feed, result: RefreshSuccess) => void, + options: RefreshBatchOptions, +): Effect.Effect => + Effect.forEach( + feeds, + (feed) => refreshOne(feed, fetchOne, applyOne, options.timeoutMs), + { concurrency: options.concurrency, discard: true }, + ) diff --git a/src/stores/feed.ts b/src/stores/feed.ts index 25f1898..57ce46e 100644 --- a/src/stores/feed.ts +++ b/src/stores/feed.ts @@ -4,6 +4,8 @@ */ import { createSignal } from "solid-js"; +import { Effect } from "effect"; +import { refreshFeedsBatch } from "../effects/feed-refresh"; import { FeedVisibility } from "../types/feed"; import type { Feed, FeedFilter, FeedSortField } from "../types/feed"; import type { Podcast } from "../types/podcast"; @@ -249,31 +251,6 @@ export function sameRefreshWindow( return fetched.every((e) => signatures.has(episodeSignature(e))); } -/** Run `fn` over every item with at most `limit` executions in flight — a - * classic worker pool. Workers pull indexes from a shared counter, so the - * first `limit` calls start immediately and each completion frees its slot - * for the next item; results are assembled in INPUT order regardless of - * completion order. A hung `fn` holds at most one slot. */ -async function mapWithConcurrency( - items: T[], - limit: number, - fn: (item: T) => Promise, -): Promise { - const results = new Array(items.length); - let nextIndex = 0; - const workers = Array.from( - { length: Math.min(limit, items.length) }, - async () => { - let i: number; - while ((i = nextIndex++) < items.length) { - results[i] = await fn(items[i]); - } - }, - ); - await Promise.all(workers); - return results; -} - function createFeedStore() { const [feeds, setFeeds] = createSignal([]); const [sources, setSources] = createSignal([ @@ -644,40 +621,40 @@ function createFeedStore() { })(), "Refreshing"); }; - /** Refresh all feeds — bounded concurrency (at most FETCH_CONCURRENCY - * in-flight requests), and each feed's refreshed episodes are applied - * AS ITS OWN FETCH LANDS (no Promise.all barrier). Per-feed apply is - * safe because applyRefreshedEpisodes keeps unchanged feeds' object - * identity and lastUpdated (union merge), so each feed's refreshed - * episodes render as its own fetch resolves — the order flapping the - * old atomic barrier existed to hide can no longer happen. */ + /** Refresh all feeds via the Effect batch program (effects/feed-refresh): + * bounded concurrency (at most FETCH_CONCURRENCY in-flight requests) + * and each feed's refreshed episodes applied AS ITS OWN FETCH LANDS + * (no barrier — the apply runs inside the feed's own fiber). Per-feed + * apply is safe because applyRefreshedEpisodes keeps unchanged feeds' + * object identity and lastUpdated (union merge), so each feed's + * refreshed episodes render as its own fetch resolves — the order + * flapping the old atomic barrier existed to hide can no longer + * happen. A failed or timed-out fetch (null episodes) leaves that + * feed untouched. */ const refreshAllFeeds = async () => { setIsLoadingFeeds(true); try { - await mapWithConcurrency( - feeds(), - FETCH_CONCURRENCY, - async (feed) => { - const { episodes, coverUrl } = await fetchEpisodes( - feed.podcast.feedUrl, - MAX_EPISODES_REFRESH, - feed.id, - ); - // A failed fetch (null) leaves that feed untouched. - if (!episodes) return; - setFeeds((prev) => { - let updated = applyRefreshedEpisodes(prev, feed.id, episodes); - if (coverUrl) { - updated = updated.map((f) => - f.id === feed.id && !f.podcast.coverUrl && coverUrl - ? { ...f, podcast: { ...f.podcast, coverUrl } } - : f, - ); - } - if (updated !== prev) scheduleSaveFeeds(); - return updated; - }); - }, + await Effect.runPromise( + refreshFeedsBatch( + feeds(), + (feed) => + fetchEpisodes(feed.podcast.feedUrl, MAX_EPISODES_REFRESH, feed.id), + (feed, { episodes, coverUrl }) => { + setFeeds((prev) => { + let updated = applyRefreshedEpisodes(prev, feed.id, episodes); + if (coverUrl) { + updated = updated.map((f) => + f.id === feed.id && !f.podcast.coverUrl && coverUrl + ? { ...f, podcast: { ...f.podcast, coverUrl } } + : f, + ); + } + if (updated !== prev) scheduleSaveFeeds(); + return updated; + }); + }, + { concurrency: FETCH_CONCURRENCY, timeoutMs: FETCH_TIMEOUT_MS }, + ), ); // Global auto-download: one idempotent pass after the batch. runAutoDownload(); diff --git a/tests/feed-refresh-effect.test.ts b/tests/feed-refresh-effect.test.ts new file mode 100644 index 0000000..2e60624 --- /dev/null +++ b/tests/feed-refresh-effect.test.ts @@ -0,0 +1,230 @@ +/** + * Feed-refresh Effect program tests (src/effects/feed-refresh.ts). + * + * These test the Effect program in isolation — no store singleton, no + * network, no fake timers. The fetch/apply closures are injected, and the + * `Clock` service comes from TestContext's TestClock, so timeouts are driven + * deterministically with TestClock.adjust instead of real 20s waits. + * + * Contracts pinned here (mirrored at the store level by + * feed-nonblocking.test.ts / feed-refresh.test.ts against a real Bun.serve): + * 1. Bounded concurrency — never more than `concurrency` fetches in + * flight, and the pool pulls the next feed as one completes. + * 2. Per-feed apply as its own fetch lands (no barrier). + * 3. A timed-out fetch leaves that feed untouched and does not stall the + * batch (TestClock.adjust fires the timeout deterministically). + * 4. A rejecting fetch leaves that feed untouched and does not fail the + * batch. + */ + +import { test, expect } from "bun:test" +import { Duration, Effect, Fiber, TestClock, TestContext } from "effect" +import { + refreshFeedsBatch, + type RefreshFetchResult, +} from "../src/effects/feed-refresh" +import type { Feed } from "../src/types/feed" +import type { Podcast } from "../src/types/podcast" +import type { Episode } from "../src/types/episode" + +const makePodcast = (id: string): Podcast => ({ + id, + title: `Show ${id}`, + description: `Show ${id} description`, + feedUrl: `http://example.com/${id}.xml`, + lastUpdated: new Date(0), + isSubscribed: true, +}) + +const makeFeed = (id: string): Feed => ({ + id, + podcast: makePodcast(id), + episodes: [], + visibility: "public" as Feed["visibility"], + sourceId: "test", + lastUpdated: new Date(0), + isPinned: false, +}) + +const makeEpisode = (id: string): Episode => ({ + id, + podcastId: "pod", + title: `Ep ${id}`, + description: "", + audioUrl: `https://example.com/${id}.mp3`, + duration: 60, + pubDate: new Date(0), +}) + +/** Resolve an episode result without dragging in the full RSS shape. */ +const ok = (episodeIds: string[]): RefreshFetchResult => ({ + episodes: episodeIds.map(makeEpisode), + coverUrl: undefined, +}) + +/** One macrotask turn — lets microtask-scheduled Effect fibers run. */ +const tick = (): Promise => { + const { promise, resolve } = Promise.withResolvers() + setImmediate(resolve) + return promise +} + +/** A resolvable fetch gate: the pool parks on `promise` until the test + * resolves it. (Promise.withResolvers's return type is not in tsconfig's + * ES2015.Promise lib, hence the explicit shape.) */ +interface Gate { + promise: Promise + resolve: (value: RefreshFetchResult) => void +} + +/** Poll `cond` across up to `iterations` event-loop turns. */ +async function pollUntil( + cond: () => boolean, + iterations = 500, +): Promise { + for (let i = 0; i < iterations; i++) { + if (cond()) return true + await tick() + } + return cond() +} + +test("bounds in-flight fetches to the configured concurrency", async () => { + const feeds = Array.from({ length: 10 }, (_, i) => makeFeed(`feed-${i}`)) + let inFlight = 0 + let maxInFlight = 0 + const gates: Gate[] = [] + const applied: string[] = [] + + const program = refreshFeedsBatch( + feeds, + (feed) => { + inFlight++ + if (inFlight > maxInFlight) maxInFlight = inFlight + const gate = Promise.withResolvers() + gates.push(gate) + return gate.promise.finally(() => { + inFlight-- + }) + }, + (feed) => { + applied.push(feed.id) + }, + { concurrency: 4, timeoutMs: 60_000 }, + ) + + // Run the batch in flight (NOT awaited) and observe the pool from + // outside via the gate side effects. + const done = Effect.runPromise(program) + // The pool starts exactly `concurrency` fetches up front. + const sawStart = await pollUntil(() => gates.length >= 4) + expect(sawStart).toBe(true) + expect(maxInFlight).toBe(4) + expect(gates.length).toBe(4) + + // Resolve one gate: the pool pulls the next feed, still bounded at 4. + gates[0].resolve(ok(["a"])) + const sawPull = await pollUntil(() => gates.length >= 5) + expect(sawPull).toBe(true) + expect(maxInFlight).toBeLessThanOrEqual(4) + + // Release everything, re-draining as the pool pulls new gates, until + // every feed has been fetched and applied. + while (applied.length < 10) { + for (const gate of gates.splice(0)) gate.resolve(ok(["x"])) + await tick() + } + await done + expect(maxInFlight).toBeLessThanOrEqual(4) + expect(applied).toHaveLength(10) +}) + +test("applies each feed as its own fetch lands (no barrier)", async () => { + const a = makeFeed("a") + const b = makeFeed("b") + const applied: string[] = [] + let bCalled = false + const gateB = Promise.withResolvers() + + const program = refreshFeedsBatch( + [a, b], + (feed) => { + if (feed.id === "a") return Promise.resolve(ok(["a-1"])) + bCalled = true + return gateB.promise + }, + (feed) => { + applied.push(feed.id) + }, + { concurrency: 4, timeoutMs: 60_000 }, + ) + + // Run the batch in flight; A's fetch resolves and applies while B's is + // still parked at the gate. + const done = Effect.runPromise(program) + const aApplied = await pollUntil(() => applied.includes("a")) + expect(aApplied).toBe(true) + expect(bCalled).toBe(true) + expect(applied).toEqual(["a"]) + expect(applied).not.toContain("b") + + gateB.resolve(ok(["b-1"])) + await done + expect(applied).toEqual(["a", "b"]) +}) + +test("a timed-out fetch leaves that feed untouched, without stalling the batch", async () => { + const fast = makeFeed("fast") + const hung = makeFeed("hung") + const applied: string[] = [] + // A promise that never settles — the fetch hangs past the timeout. + const never = new Promise(() => {}) + + const program = refreshFeedsBatch( + [fast, hung], + (feed) => + feed.id === "fast" + ? Promise.resolve(ok(["f-1"])) + : never, + (feed) => { + applied.push(feed.id) + }, + { concurrency: 4, timeoutMs: 5_000 }, + ) + + const timed = Effect.gen(function* () { + const fiber = yield* Effect.fork(program) + // Advance the TestClock past the timeout: the hung fetch's + // Effect.timeout fires deterministically — no real 5s wait. + yield* TestClock.adjust(Duration.millis(5_000)) + yield* Fiber.join(fiber) + }) + await Effect.runPromise( + timed.pipe(Effect.provide(TestContext.TestContext)), + ) + + // The fast feed applied; the hung one was dropped, and the batch + // completed anyway. + expect(applied).toEqual(["fast"]) +}) + +test("a rejecting fetch leaves that feed untouched and does not fail the batch", async () => { + const bad = makeFeed("bad") + const good = makeFeed("good") + const applied: string[] = [] + + const program = refreshFeedsBatch( + [bad, good], + (feed) => + feed.id === "bad" + ? Promise.reject(new Error("feed exploded")) + : Promise.resolve(ok(["g-1"])), + (feed) => { + applied.push(feed.id) + }, + { concurrency: 4, timeoutMs: 60_000 }, + ) + + await Effect.runPromise(program) + expect(applied).toEqual(["good"]) +})