From 69cebb5d15939bf9b6c1b4647b53879beab91ba2 Mon Sep 17 00:00:00 2001 From: Owen Lin Date: Wed, 2 Sep 2026 21:21:54 +0000 Subject: [PATCH] Route rollout reads through the canonical JSON decoder (#42378) ## Why Directly deserializing the flattened `RolloutLine` envelope can reject nested decimal values, preventing affected paginated sessions from resuming. ## What changed - Add canonical string, byte, and reverse-scanner helpers that decode rollout records through `serde_json::Value` before decoding the flattened item. - Route rollout readers across session discovery, history, migration, search, thread storage, and transcript previews through those helpers. - Remove `Deserialize` from `RolloutLine` so new readers cannot bypass the canonical persistence decoder. ## Testing Add coverage that resumes a paginated rollout after a token-count record with a decimal rate-limit value and verifies that ordinal sequencing continues. GitOrigin-RevId: 49abac1e0751c073daa5a93a840d8a483fd2d013 --- .../app-server-exports-stable.json.zst | Bin 147122 -> 147183 bytes .../tests/suite/v2/git_attribution.rs | 3 +- .../app-server/tests/suite/v2/thread_fork.rs | 2 +- .../app-server/tests/suite/v2/thread_list.rs | 3 +- .../tests/suite/v2/thread_revert.rs | 3 +- codex-rs/cli/src/doctor/thread_inventory.rs | 14 ++-- codex-rs/core/src/agent/control_tests.rs | 5 +- codex-rs/core/tests/suite/abort_tasks.rs | 3 +- codex-rs/core/tests/suite/agents_md.rs | 3 +- .../tests/suite/collaboration_instructions.rs | 3 +- codex-rs/core/tests/suite/compact.rs | 11 ++- codex-rs/core/tests/suite/compact_remote.rs | 11 ++- .../core/tests/suite/compact_remote_parity.rs | 3 +- .../core/tests/suite/compact_resume_fork.rs | 5 +- codex-rs/core/tests/suite/fork_thread.rs | 3 +- codex-rs/core/tests/suite/guardian_review.rs | 3 +- codex-rs/core/tests/suite/hooks.rs | 3 +- codex-rs/core/tests/suite/image_rollout.rs | 7 +- codex-rs/core/tests/suite/items.rs | 3 +- .../core/tests/suite/mcp_tool_exposure.rs | 5 +- .../core/tests/suite/mcp_turn_metadata.rs | 3 +- codex-rs/core/tests/suite/model_switching.rs | 3 +- codex-rs/core/tests/suite/pending_input.rs | 3 +- codex-rs/core/tests/suite/remote_env.rs | 5 +- codex-rs/core/tests/suite/review.rs | 15 ++-- codex-rs/core/tests/suite/settings_commits.rs | 2 +- .../tests/suite/spawn_agent_description.rs | 3 +- .../core/tests/suite/token_usage_rollout.rs | 3 +- codex-rs/exec/src/lib.rs | 2 +- codex-rs/history/src/lib.rs | 6 +- codex-rs/history/src/tests.rs | 14 ++-- codex-rs/rollout/src/lib.rs | 14 +++- codex-rs/rollout/src/list.rs | 7 +- codex-rs/rollout/src/ordinal.rs | 5 +- codex-rs/rollout/src/recorder_tests.rs | 71 +++++++++++++++++- codex-rs/rollout/src/reverse_jsonl_scanner.rs | 14 ++++ codex-rs/rollout/src/search.rs | 3 +- codex-rs/rollout/src/tests.rs | 2 +- .../thread-store/src/local/model_context.rs | 3 +- .../src/local/model_context_tests.rs | 4 +- .../src/local/rollout_lineage_tests.rs | 2 +- .../local/rollout_migration/line_parser.rs | 2 +- .../rollout_migration/line_parser_tests.rs | 4 +- .../src/local/rollout_migration/publish.rs | 3 +- .../src/local/rollout_migration_tests.rs | 2 +- .../src/local/thread_history/read_tests.rs | 2 +- .../thread_history_materialization_tests.rs | 4 +- .../src/resume_picker_transcript_preview.rs | 3 +- .../resume_picker_transcript_preview_tests.rs | 1 + 49 files changed, 184 insertions(+), 114 deletions(-) diff --git a/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst b/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst index 719e533b6ab1fed3ffb0face4e40f9f3a05ab894..e499e34c1fbe2d3b0a09756ee55ae4152f9d7130 100644 GIT binary patch delta 4435 zcmV-Z5v=aA{Rr><2!MnEv;xd6e+7vE=GR9bAd?$#*fdQ&MS}$8yAw}yi!>&)KOauLXkA(zMiO%=6=bOb-oa%Z(vNHU>lrY0wYV|I9>kiATf8MD{Xu+UE znL&pcRG7hyW))^xs##ZGzFYy%Zqp2>dN7Tgyjgcri~*9xnfk`bcCUBQdh0`SaV|H_ z&p(;5?z;b|-6)E7ZNZQgpPoqOxY^qcK^v-n6jInqhVpl-Rxqa&102bmttJ4 zC81{NwEUV4m-c+Q?C9fpfANq&AS`}|ObG#vkH=6Dh|#%8@bBzYy9S9gG}-z%nHVuI z^AfdNg&e4}9hXUxX%15$5J*sv01=5m6pzH>aakp%`2!Qc!9;_BF%&Waf&d{H1Rw|o z0*)XEf*=S1iV!3Si!dZYR3>|i0Zc#X$9BZiX~f9Pfz@qe+%5v)q< zb2NK!oQ#o@shXU*4AFz1`3vU6QvE$R;=JFqZW082FZvH!K7J{u%y}oGCckzbLz!GM z4RjuqnTU<`E?P9&r?*TEI!#TsNSYU6!~r(pFG&;)Y$e%a#aqS4Rd3yaUsw&;wRm4C zIX2p}xaS*mC78Fyf0fmW4VLRxy?zb$b8F3>bhNdRpRkqvufh-gZL9OT5y-E<&i*uf ztiF01{ltqDwa73cE1OazozNgyh?_- zOAo}|L9;s_W~LZE_xBf5)9yFK@;t5;8SI0gFei-XM$mY%e>;({w1jjrGG}rwQ%Cay z-Z2AJ@V-1a4UuE$tM_pdIszKYM2AytDYMDZY2P2Y`KIwe?j3{;NPz4fkhvJh5Uv#% zQg)=2pFmZ=5Yz<#CrFmhti52gA3M0Pn9DP(+5HF*wAhKrWFQ7Rs z`Kgmp#W@F%02K`g2X21x$Wr~hR$cOr8<|!>$~8*`f3I3j%r<0G_M>?KW#1Z&$(Xj$ zqfs^sTZIq6IIm2rkwIHUZh=Fy9j4T@uSr)3vle@oE9CmhN^8cHJ9+cALbvcXjt)jK zs!#wr6S_xq!%^1rV$4F&1-ivX0ipbb+^vaW-xNn)R#GS&hi5f{sw04vP-4F-=st~9 zo#L6He^gl9uOkW7CXO4x5x`8;#eVI5NdE!DC>z^LfXr6iAZI3n$=*8kX)Yogi3S6X zK+Y9(?j?5O6s^b)nom2K^OaGnJ2ypsL<#QtN;ERVef>k6DnT`Y>ue`*4pi{MQ|Tyl z3!En`pim2T#n#uiB-G4;cJy)`p$vnqx4&)ff2%w=*tFX-oM1u>+Ox*h5%fPT-KJ|k zjUFg7t&3h~wC@!r`KC**-uNGS(8TMb(74X-mq6#y)~v7YREKLD&X+;)Ygbp-;-xDE zlh4gMd_w1;h3lu|*61<1Gj5p<0`y2}P4G|Y`CK%`*JpayY6L^eaK55)rveq5m0k{W ze+uq~LY->ga1mpCug|^QSQ`n*8A75PL|A5QFqD5%Y;aPC2gn+w;Ektk*l@!~;?N96 z&evg8!To7CAybvwAkwx5vqJJnm0nvesGsAnGy&BqXc6He-czC?{?s4ct@%WSaZ_^i zrp|Wp(|Vm6*@&-`>7$GWQbyC&NBQdxe`3`J*he|)Qq_XzI@A&r6tM|C%5Y)B1MY7d zrk7d|z`X<|3)tnkdIsK3F0D%Cim z8Ve*)N;8uO=ndGVi^**5g-oZu)5YWnJij=sdA2}8NfyK3FzEu%4lXdC%H1mzf3UX4 zzWw6{#$UA7Qf5@hj2n&BQ6ZDuyJG~5N@`q5>eQ4si>(u0JJK~ScE<}%+z`MN1WB3@ z_$2}=?l!p%oXll~_c*A75LZ0GzTDSe%>Olqd`2#W7X3$T?wYLWJHDh>a-EuIEGWl+ zo+!Gd69*SIpfZpytl^Y91)GgVf9nSMZnTlPq|6WQSg>13R9?=)IhghwElU}n|K=9A zp^+C=2E9XT+L^8GR2L$GT!0xWD_REM1n)3o_y_WMh`Jb1{A{i2Z&DMhU+z(amvK$c(f7NK@FKP87 zOjK9`Fzv@}7Ue$g>q`u2&++Yr4>yroBVAj?zTl!3O1BS*x4pSB9fH`$K%awJ2 za}1>SH{$7~rw(gffv^8jtX;Sk`mCgN$qG+tpPChG;6TZR#6%x>;%E*~GOdQ~J4%}1 z0eDq@goz}Aczp22k29$0f1kYyGoS|4!AlnTh9=$jV`<6WMJc69wiPm^8#*O>!ON26 zGHH(ySXMsrW@XiK7!t&(J@)Z08i8hpm*V*`!OPhzyP(DAnGOL6s&{)AG;9v&l;@Oo z%GeVdA9de>#1h@;k$vskdpwUStijeI_elB~vg!WxUVmeq))gg0wBly0-M|tQWCSQg zaL1guXLV=lIDUx0e?WUO_%t{Xgi2udCxtkkLvDVDc)sIOLu70e_ z(@Pzz`jXA6qZjra7$IqWfC5(m%*~P0qG4+Mv91^+4 zLfwv`7#GSOva$DJm3mvSL;^+if2cgOCpSA==&%;gPZq%p)o_Qr z6;Md00y5|BQgfK!MpZywErW?t3z@nmR~;Lmmuewz7fWNrd3KF3{X^1J?*HKn z$qIXfNskElx7(mzHi~sX3#Yzn_|FEnJe|Yo8@u)xE_@Z?VC5RVhqKs#&SHVGT5@i@ z*jnbbf1%btEwid2I)W460O5~7&OFnw`}}%FbPZ-{RFRL~SW9qojFhG^{Gq4W;xR~5 zQe0!_FPk%cTR*eWqVUR@zr*LNWka8aEDQz6zAJY}p^rBt6@)@{;DBbFMQ3Jb8(jgv zW%o+|<%in+HyivJEh0Si)Bw=bepOiaF)Br&e=zZ+@d46*cbk{3nSboH?nptY#-0>m zl?EA(bp$X_8D&&9O}o;E*U)bOl1w2Hf!Tn|1PMvnZhYJo3z>s}@mpBU2b$N#HTWSv+B* z_2U2j)8+UNs@R>%x(icWw<}F0iVHQrb#1cWSqw3w3fRjQ;liRinX=&>Ncy8M76Fc6 zi``-Dir<Bt}In&vj2g(vB3`5W(9-pY^t<5nxSp7Q?A&iC>70e>SXRN)1v~MO|e?^B2gI< zB3}hbq1Y7>i@*^s2!t`HjFM;^$Lx6aC*qSp4!ZU$H(`@10L4;wlf)#5T_t0*@EsFB7*ou)L&EK*I?%m}+9_2nv0$A+O-OOaDL zwiI>i!$=s80_OTLa56$sn z5i1c69j(7ukC!>K|VRG)9JGu zRzkA3^8l8!0-W0X6TIyaf^p*N?gYJ+uh!@6W`l)>SkKpEj1Um{K zxE!EcT-}TsVj)9e_#NZ+Q|yT>*K$xWuYQ5RfN7B*@miN0ir#gUJQx8BN%?1Eg;Mqw zRGh#=a@%*3n<}xg)S6)ezq;B%0Qx};J4Wnyf_k?ZZNLVce{4m=NgwFQD7HB&idvEm zQ|sZM-}e>~{+T(v;!~nF;NE^NV@+dPh~7<~WG^HTwl=9DIjjQz^v_=Gj+UNPv1)%o z9uQC|pLf)8*mc7n9949Yj?NfV8bXTS=28IN!>t_WNl}b*TPaOnUPl0G1wspa|GtcZ zv&?-0PeeGoe=S1N?hv@{`RRu%Js`9AF?=khOMLtQ4e?oCVk5M}D6IG~EHZJM`h{uif4X|PFFs#zV~Ksuy23=l-@?d@eS31JdjDwWa9BR{XHZZ@v=lFdTfp$v)yWP@mJv9tYyvJ~l3dp-fE*bB=D?V?Ot5?tl>J+A$oV(W z-wP)TMP-v!pP!QoLs@L=fTLq{7bGfs&|y0)^c?&gF&-;SH%)JbGLWRJ)gOwR_ ZI)m@{AJJx8MBAt_>{XC1uXy)so&(*|f}Q{X delta 4374 zcmV+x5$W#l{Rp!C2!MnEv;xd6e;tPa1{X#kFq0dg*^@m*N(hYhL+f>)So>WD%bcZp zSNE?YA|jvyJODcYI{?tJ=}`UNX)uU9tGJ%qk?*wGF;+nT&dQ}GAGANCXP?gZA0?v% zZqq1N_1ty}xQ~;^LISBo=lj|7r7@E-{QjJfVm>AH4tnbjqg8LsCbVGCf8nIiA%zlB zXh}+8RV6j+(w8rnfM;hl!zw$NW>&V=T@)kB`U3Bu^_CCG#rd?FpT9C=-R1t9R--7| zvjsy|Y(0_G${hEqmfIvr?EaOxP-izkN5jI9yV*r4R_S^@0Uzrzsy4r7!<9W>J{#IP zFCG#IM8)qwB_SZ=<1rKje`0h_68t;+)Sf{ijcm4jn@o(Dm-%Q~EvlGG*EV$b3WW}+ z(+!nLl4%Z8AP`7UfB+GRKopR~!f|mWrhOBD!9;^`F%&Waf&d{H1Rw|of{q{vf*=S1 ziV!3Si!dNUWTwiC0l+@#$9B1WQn)(p%1yyW=9@o8=pHk_e_{nAf0#u3=E(T?AsJ&N zYf8@2_0QKovmAzHsm_1+ad>+uZlDF7nd^r&AKxG+&v~D-Care5eoSf)4f!6hnTW~m z7B7wVs@skRopFtvCGGqW`T_Q}QKZBM_Dfd3;%)e&dxz{e&!+}FT6`sg-Wu&G-2Vn$ zam*9rstv&g__|{xe-gfSa9d|>BJES8A?!oUt++vMuIljD>+PR5#QH7icQPuOvk4oS z5gO96Mz0a0c~nrPVDagZ5Xk?}J%9n*O67)=p!w4i`BO~!&l63Xpw(y~l{-I9r6f4_ z(ih{RX>Ew{KeXi*@*Yr;C+zqm&_c202v54pC>cKoZZ47$fAoNjH4~fXzUnyI=k@vC zO`6n>PzLJveJOpX?9I5c?{+ADwlZR31~~>U;g^P@D46vE6ZYd#wm90<02VpEYJ{=b zwb);mI@}hYl|*R(6%_KfwZcQmJ^jC!W2IA>B+ zqoRS!ft#;lLF6hf-lCPJ)2n_L7!yl=2vrj^gy4|YV zD_lQWzD9T0ohy~;!0jW2ZjwKx=bJ@Se7!omRwFQ12J=-&cPiXtvu1LbvwAo5q*VK6 z#Teszef5_bwULsXAwqP6ypb8YLwS3O4SM77fI_3>xbXunHk{RwID*5-DGRJB=>9aE zf2dHE4n*4Y!K|R)Ql%9w7pQSyDNTSb1uY^(!h4FZh@Ym9?$)#DGj2+Z-qhJj__V%W zBO8&%$@EdSqzHbx`Y1o%K_IyS_EFYds+#d!_bov|pEjXK*;Y0@pz*e0`iy$O#7p?U zfL(r2&%iZgdR-Dd%OIQ*lL@FCWk|*?e*y0>4-xB#yy#IU-_9Tf_|3Gfp|MK>r2=Ns z0lfjyx|odCUdUGaSWpiBc_R20atS2gTef?C*3!BX~`m^^j0NX51E>BPs#4%Wy#8&v_}D!mB)Fr zvT!+%3jov}`}o=nLOg?`e|UZh;^pX-UCGWjdvuGWCSUN8Jk` zt|T{lSYG?Ko}5QTWw4FNJ@i#xvgy9K*YDD7$x>b}$#A*3n-j5*+^K zI?Bfy!AOfigRsRcJ^vTIzL2etKN7zh4>Iqx(9uN{6(2$~-|FLIe_!t0R>hRj*T)7Q z9pEvJmM=p-E+&p4!n74T+^c?$Ofa05s2-XJluJbc%)7b0V)u zTZXO_ODAM59DECmK^;^ul7feN*bPpixC8fq3$C$gINrnDaC|NxB|ckD1uW2cA+p~oZ>(l)61Qt}o9r91WM>-Xd4enm; zBW*US0&=M=OO#s3{%g9)!PC7|3;Bvz8l&d3uiI1^e~CmC<$G7M{tsV#R@fty9&_-I zY=e5qRjdhGIRC2QayGc-lf|ar;I&76hc_SxD_8y=PL~CAUk1whmUAO%Ybmo0{r|K~ zw}$8l>;VS|IRatknT84H4`4)B8KhB_VD!dX;mxtHX&P=0oq1Y33^(${HL8Bukm=i) zncX)Ef1k|!9dchS4SkYtVXT4dyPl4wdAu=J5DL}j0-Cj0bY?o+ND63?-G}eLKWg`n z*5J<&DB%gF2AEIn&u{BKa48B99uJGpCH;5Nyi7IoZ>@C~CrXv(qyVcl$go&Pz<|mq zBa&;{iG6s%eFIQw3b8C^1N0N*JZXF8YggK8f42zp*c<%)tS6asaZi4NsZCA3p0$PE zPtPAu#7AG%RT<0OmMk;4)DY~VgyZzv-|9>~7&LA#70dh&P|Kp8(!*_*fVm@6m8bYI zs%ob=^(N)1e^0P>QY}XcHa)XhL;B#6r%bw*5zvk7aNx=dk73bbyfIoYq3l;(jz_tQ zf89x}y8y*?Bl%RK08#TRXp`NY#V8$BIK6BM78b45l=Zj+DMI>s5#R{6*!_`Rk*sNI zs6D|bF49~q3lis(Ee>(|)UcCpFZIkP`VlKrte>6+Alet@2j|%)~3Me9G;$I;5e^-#(^Zd{ zUG$Ry01nEEWA2hrl?w#0_Bvhkf22LtZaeCbwc&8cO85Y5Mj?PMHC>+E)~zGSC?S(8 zw5_|LA6uS6wsf#0MGul%%&s&b zYvnA%NG)tExgEi z`PDWz4}OBy&{AbJ3{m!z!g%~m2R8-+ndO;JgLg-w6X)T$2?=gax2n>eOD&pnb5R1@ z)Y`;Tb?td0k>zYF$52V2f0&xUqU|RF;PT=lXotYDeOFQK_d`*&Uy?C!X6^$8aB< zUOk@32{#?1gx>oz-ihe!X=O^x@*2Ne;bXR~1`8!)|1?29xvbMMe`YaWa%69H04#?p zIR4cqcr73#!^G8@1-)IV*nB^3hPG{I2l7Ay8&C_ssh*z{@}@Q0`JU9hxpO(x7IeNi z=++ZwVG927Q;UaA->HT@K<15153GUWWT|)i+`Sd}+v4BM57Rb0H!KJE)8`jUf=dBc9m2kbuS7%tLax9DdVZgxDQXe|lDwYCl~c5I8BHY81I60C^mkU7Nv4Yq zxlqVinq21|f9HK^P!3kAS)Fwt-Bx@M5JYU}+asHVYmRk@h0J-}Y0|ZScSD@l^%MHe zu9P#|=qG6M%tOd6!azqnR?2uyrmP#ig&Q*LuI=Im4rrMkgqv!&p}_iqtW_N6p+0GP zYMYaeNmSlHmT)*MbNn?^Q2n$NNuyloqUm03!a;F}f9tY<9O!qIrM5n8B|)-cSS&Z( zxFljns%h}yWCt(H2$rnu4Hxk+T-OCarV#+s#+X(Nw)|_#{;kND^Z(J`vzLXE$?}w+ zlL%7TwspW=%DW2{l|3k9J273Splf{#HC6BWg2G{5h5(_`z&nGL8FHQJcl@`lu`S|Z Q)adP12#}hqc(FYE1NvZSX8-^I diff --git a/codex-rs/app-server/tests/suite/v2/git_attribution.rs b/codex-rs/app-server/tests/suite/v2/git_attribution.rs index c96d230880..5925dd3d4f 100644 --- a/codex-rs/app-server/tests/suite/v2/git_attribution.rs +++ b/codex-rs/app-server/tests/suite/v2/git_attribution.rs @@ -27,7 +27,6 @@ use codex_config::types::AuthCredentialsStoreMode; use codex_protocol::models::ContentItem; use codex_protocol::models::ResponseItem; use codex_rollout::RolloutItem; -use codex_rollout::RolloutLine; use core_test_support::responses; use core_test_support::skip_if_no_network; use pretty_assertions::assert_eq; @@ -334,7 +333,7 @@ fn replace_attribution_fragment_with_legacy( .lines() .filter(|line| !line.trim().is_empty()) .map(|line| { - let mut line = serde_json::from_str::(line)?; + let mut line = codex_rollout::parse_rollout_line(line)?; if let RolloutItem::ResponseItem(response_item) = &mut line.item && let ResponseItem::Message { role, content, .. } = &mut response_item.item && role == "developer" diff --git a/codex-rs/app-server/tests/suite/v2/thread_fork.rs b/codex-rs/app-server/tests/suite/v2/thread_fork.rs index 830d7e713d..ee7b6fa311 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_fork.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_fork.rs @@ -1778,7 +1778,7 @@ async fn assert_thread_fork_freezes_active_paginated_turn_as_interrupted( .expect("fork history base"); let child_rollout = std::fs::read_to_string(forked_path.as_path())? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()?; assert!(matches!( child_rollout.as_slice(), diff --git a/codex-rs/app-server/tests/suite/v2/thread_list.rs b/codex-rs/app-server/tests/suite/v2/thread_list.rs index 76b2741d1c..fb8b590b4e 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_list.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_list.rs @@ -44,7 +44,6 @@ use codex_protocol::protocol::MultiAgentVersion; use codex_protocol::protocol::SessionSource as CoreSessionSource; use codex_protocol::protocol::SubAgentSource; use codex_rollout::RolloutItem; -use codex_rollout::RolloutLine; use codex_rollout::append_rollout_item_to_path; use codex_rollout::read_session_meta_line; use codex_state::DirectionalThreadSpawnEdgeStatus; @@ -218,7 +217,7 @@ fn set_rollout_cwd(path: &Path, cwd: &Path) -> Result<()> { let first_line = lines .first_mut() .ok_or_else(|| anyhow::anyhow!("rollout at {} is empty", path.display()))?; - let mut rollout_line: RolloutLine = serde_json::from_str(first_line)?; + let mut rollout_line = codex_rollout::parse_rollout_line(first_line)?; let RolloutItem::SessionMeta(mut session_meta_line) = rollout_line.item else { return Err(anyhow::anyhow!( "rollout at {} does not start with session metadata", diff --git a/codex-rs/app-server/tests/suite/v2/thread_revert.rs b/codex-rs/app-server/tests/suite/v2/thread_revert.rs index 0bb3adf4ac..ba61faf83c 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_revert.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_revert.rs @@ -38,7 +38,6 @@ use codex_protocol::config_types::Settings; use codex_protocol::openai_models::ReasoningEffort; use codex_protocol::protocol::EventMsg; use codex_rollout::RolloutItem; -use codex_rollout::RolloutLine; use codex_rollout::read_session_meta_line; use codex_utils_absolute_path::AbsolutePathBuf; use pretty_assertions::assert_eq; @@ -108,7 +107,7 @@ async fn thread_revert_preserves_fork_cutoff_after_cold_resume() -> Result<()> { let inherited_revert_cutoff = std::fs::read_to_string(parent.path.as_ref().expect("parent rollout"))? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()? .into_iter() .find_map(|line| match line.item { diff --git a/codex-rs/cli/src/doctor/thread_inventory.rs b/codex-rs/cli/src/doctor/thread_inventory.rs index 702195be64..971c32f831 100644 --- a/codex-rs/cli/src/doctor/thread_inventory.rs +++ b/codex-rs/cli/src/doctor/thread_inventory.rs @@ -5,7 +5,6 @@ use super::Config; use super::DoctorCheck; use super::DoctorIssue; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_protocol::protocol::InternalSessionSource; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SubAgentSource; @@ -543,7 +542,7 @@ async fn thread_id_from_rollout(path: &Path) -> RolloutThreadId { Err(_) => continue, }; if item_type == "session_meta" { - return match serde_json::from_str::(line.trim()) { + return match codex_rollout::parse_rollout_line(line.trim()) { Ok(line) => match line.item { RolloutItem::SessionMeta(session_meta) => { RolloutThreadId::Id(session_meta.meta.id.to_string()) @@ -560,7 +559,7 @@ async fn thread_id_from_rollout(path: &Path) -> RolloutThreadId { }; } if !has_legacy_item { - has_legacy_item = serde_json::from_str::(line.trim()).is_ok(); + has_legacy_item = codex_rollout::parse_rollout_line(line.trim()).is_ok(); } } @@ -745,6 +744,7 @@ where #[cfg(test)] mod tests { use super::*; + use codex_history::RolloutLine; use codex_protocol::ThreadId; use codex_utils_absolute_path::test_support::PathExt; use pretty_assertions::assert_eq; @@ -862,7 +862,7 @@ mod tests { fixture.write_rollout(/*archived*/ false, "2025-01-02T10-00-00", filename_id); let contents = std::fs::read_to_string(&path).expect("rollout file"); let mut rollout_line = - serde_json::from_str::(contents.trim()).expect("rollout line"); + codex_rollout::parse_rollout_line(contents.trim()).expect("rollout line"); let RolloutItem::SessionMeta(session_meta) = &mut rollout_line.item else { panic!("expected session metadata"); }; @@ -981,7 +981,7 @@ mod tests { fixture.write_rollout(/*archived*/ false, "2025-01-02T10-00-00", filename_id); let contents = std::fs::read_to_string(&metadata_path).expect("rollout file"); let mut rollout_line = - serde_json::from_str::(contents.trim()).expect("rollout line"); + codex_rollout::parse_rollout_line(contents.trim()).expect("rollout line"); let RolloutItem::SessionMeta(session_meta) = &mut rollout_line.item else { panic!("expected session metadata"); }; @@ -1114,7 +1114,7 @@ mod tests { fixture.write_rollout(/*archived*/ false, "2025-01-02T10-00-00", filename_id); let contents = std::fs::read_to_string(&path).expect("rollout file"); let mut rollout_line = - serde_json::from_str::(contents.trim()).expect("rollout line"); + codex_rollout::parse_rollout_line(contents.trim()).expect("rollout line"); let RolloutItem::SessionMeta(session_meta) = &mut rollout_line.item else { panic!("expected session metadata"); }; @@ -1148,7 +1148,7 @@ mod tests { fixture.write_rollout(/*archived*/ false, "2025-01-02T10-00-00", filename_id); let contents = std::fs::read_to_string(&path).expect("rollout file"); let mut rollout_line = - serde_json::from_str::(contents.trim()).expect("rollout line"); + codex_rollout::parse_rollout_line(contents.trim()).expect("rollout line"); let RolloutItem::SessionMeta(session_meta) = &mut rollout_line.item else { panic!("expected session metadata"); }; diff --git a/codex-rs/core/src/agent/control_tests.rs b/codex-rs/core/src/agent/control_tests.rs index 0e787e122c..a6ba2c49d1 100644 --- a/codex-rs/core/src/agent/control_tests.rs +++ b/codex-rs/core/src/agent/control_tests.rs @@ -22,7 +22,6 @@ use codex_extension_api::empty_extension_registry; use codex_features::Feature; use codex_history::CompactedItem; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::AuthManager; use codex_login::CodexAuth; use codex_protocol::AgentPath; @@ -1262,7 +1261,7 @@ async fn spawn_agent_fork_from_paginated_parent_uses_model_context_prefix() { let lines = std::fs::read_to_string(&rollout_path) .expect("read child rollout") .lines() - .map(|line| serde_json::from_str::(line).expect("parse rollout line")) + .map(|line| codex_rollout::parse_rollout_line(line).expect("parse rollout line")) .collect::>(); let RolloutItem::SessionMeta(meta_line) = &lines[0].item else { panic!("child rollout should start with session metadata"); @@ -1466,7 +1465,7 @@ async fn spawn_agent_fork_drops_inherited_token_usage_state() { let lines = std::fs::read_to_string(&rollout_path) .expect("read child rollout") .lines() - .map(|line| serde_json::from_str::(line).expect("parse rollout line")) + .map(|line| codex_rollout::parse_rollout_line(line).expect("parse rollout line")) .collect::>(); assert!( !lines.iter().any(|line| { diff --git a/codex-rs/core/tests/suite/abort_tasks.rs b/codex-rs/core/tests/suite/abort_tasks.rs index 4b995595d5..797ba476fe 100644 --- a/codex-rs/core/tests/suite/abort_tasks.rs +++ b/codex-rs/core/tests/suite/abort_tasks.rs @@ -3,7 +3,6 @@ use codex_core::StartThreadOptions; use codex_core::SuspendTurnOutcome; use codex_core::TurnInputRequest; use codex_history::RolloutItem; -use codex_history::RolloutLine; use std::sync::Arc; use std::time::Duration; @@ -152,7 +151,7 @@ async fn root_turn_suspension_preserves_unfinished_turn_history() { let items = rollout .lines() .map(|line| { - serde_json::from_str::(line) + codex_rollout::parse_rollout_line(line) .expect("parse durable rollout") .item }) diff --git a/codex-rs/core/tests/suite/agents_md.rs b/codex-rs/core/tests/suite/agents_md.rs index 297bc1d60f..4a93b247f4 100644 --- a/codex-rs/core/tests/suite/agents_md.rs +++ b/codex-rs/core/tests/suite/agents_md.rs @@ -8,7 +8,6 @@ use codex_exec_server::LOCAL_ENVIRONMENT_ID; use codex_exec_server::REMOTE_ENVIRONMENT_ID; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_home::CodexHomeUserInstructionsProvider; use codex_protocol::config_types::TrustLevel; use codex_protocol::models::PermissionProfile; @@ -98,7 +97,7 @@ fn remove_agents_md_world_state_section(rollout_path: &Path) -> Result<()> { let mut removed_section = false; let retained = rollout .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()? .into_iter() .map(|mut line| { diff --git a/codex-rs/core/tests/suite/collaboration_instructions.rs b/codex-rs/core/tests/suite/collaboration_instructions.rs index 512d4be8b6..f6a8fc3c05 100644 --- a/codex-rs/core/tests/suite/collaboration_instructions.rs +++ b/codex-rs/core/tests/suite/collaboration_instructions.rs @@ -2,7 +2,6 @@ use anyhow::Result; use codex_core::TurnInputRequest; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_models_manager::model_info::model_info_from_slug; use codex_protocol::config_types::CollaborationMode; @@ -1037,7 +1036,7 @@ async fn cold_resume_refreshes_legacy_collaboration_snapshot_once( let legacy_rollout = std::fs::read_to_string(&rollout_path)? .lines() .map(|original_line| { - let mut line = serde_json::from_str::(original_line)?; + let mut line = codex_rollout::parse_rollout_line(original_line)?; if let RolloutItem::WorldState(world_state) = &mut line.item && let Some(snapshot) = world_state.state.get_mut("collaboration_mode") { diff --git a/codex-rs/core/tests/suite/compact.rs b/codex-rs/core/tests/suite/compact.rs index 3473e39628..09ae3e2cdd 100644 --- a/codex-rs/core/tests/suite/compact.rs +++ b/codex-rs/core/tests/suite/compact.rs @@ -6,7 +6,6 @@ use codex_core::compact::SUMMARY_PREFIX; use codex_core::config::Config; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_model_provider_info::ModelProviderInfo; use codex_model_provider_info::built_in_model_providers; @@ -367,7 +366,7 @@ fn replacement_history_from_rollout(path: &Path) -> Result> { .map(str::trim) .filter(|line| !line.is_empty()) { - let entry: RolloutLine = serde_json::from_str(line)?; + let entry = codex_rollout::parse_rollout_line(line)?; if let RolloutItem::Compacted(compacted) = entry.item && let Some(items) = compacted.replacement_history { @@ -723,7 +722,7 @@ async fn summarize_context_three_requests_and_instructions( if trimmed.is_empty() { continue; } - let Ok(entry): Result = serde_json::from_str(trimmed) else { + let Ok(entry) = codex_rollout::parse_rollout_line(trimmed) else { continue; }; match entry.item { @@ -1049,7 +1048,7 @@ async fn manual_compact_records_durable_and_local_token_usage() { let rollout_items = fs::read_to_string(rollout_path) .expect("read rollout") .lines() - .filter_map(|line| serde_json::from_str::(line).ok()) + .filter_map(|line| codex_rollout::parse_rollout_line(line).ok()) .map(|line| line.item) .collect::>(); let records = rollout_items @@ -3430,7 +3429,7 @@ async fn pre_sampling_compact_recovers_comp_hash_after_resume() { let rollout = fs::read_to_string(&rollout_path).expect("read rollout"); let persisted_comp_hash = rollout .lines() - .filter_map(|line| serde_json::from_str::(line).ok()) + .filter_map(|line| codex_rollout::parse_rollout_line(line).ok()) .find_map(|line| match line.item { RolloutItem::TurnContext(context) => context.comp_hash, _ => None, @@ -3713,7 +3712,7 @@ async fn auto_compact_persists_rollout_entries() { if trimmed.is_empty() { continue; } - let Ok(entry): Result = serde_json::from_str(trimmed) else { + let Ok(entry) = codex_rollout::parse_rollout_line(trimmed) else { continue; }; match entry.item { diff --git a/codex-rs/core/tests/suite/compact_remote.rs b/codex-rs/core/tests/suite/compact_remote.rs index dd905d4f67..cfaf139ddd 100644 --- a/codex-rs/core/tests/suite/compact_remote.rs +++ b/codex-rs/core/tests/suite/compact_remote.rs @@ -16,7 +16,6 @@ use codex_features::Feature; use codex_history::CodexHarnessMetadata; use codex_history::InitialHistory; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_login::auth::AgentIdentityAuth; use codex_login::auth::AgentIdentityAuthRecord; @@ -404,7 +403,7 @@ fn annotate_retained_user_in_rollout(path: &Path, retained_text: &str) -> Result let mut rollout = fs::read_to_string(path)? .lines() .filter(|line| !line.trim().is_empty()) - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()?; rollout .iter_mut() @@ -432,7 +431,7 @@ fn assert_compacted_user_metadata(path: &Path, retained_text: &str) -> Result<() let replacement_history = fs::read_to_string(path)? .lines() .filter(|line| !line.trim().is_empty()) - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()? .into_iter() .rev() @@ -669,7 +668,7 @@ async fn remote_compact_v2_retains_only_client_developer_messages_when_enabled( codex.shutdown_and_wait().await?; let replacement_history = fs::read_to_string(&rollout_path)? .lines() - .filter_map(|line| serde_json::from_str::(line).ok()) + .filter_map(|line| codex_rollout::parse_rollout_line(line).ok()) .filter_map(|line| match line.item { RolloutItem::Compacted(compacted) => compacted.replacement_history, _ => None, @@ -738,7 +737,7 @@ async fn remote_compact_v2_records_usage_before_output_validation() -> Result<() let record = fs::read_to_string(&rollout_path)? .lines() - .filter_map(|line| serde_json::from_str::(line).ok()) + .filter_map(|line| codex_rollout::parse_rollout_line(line).ok()) .filter_map(|line| match line.item { RolloutItem::TokenUsageRecord(record) => Some(record), _ => None, @@ -3349,7 +3348,7 @@ async fn remote_compact_persists_replacement_history_in_rollout() -> Result<()> .map(str::trim) .filter(|l| !l.is_empty()) { - let Ok(entry) = serde_json::from_str::(line) else { + let Ok(entry) = codex_rollout::parse_rollout_line(line) else { continue; }; if let RolloutItem::Compacted(compacted) = entry.item diff --git a/codex-rs/core/tests/suite/compact_remote_parity.rs b/codex-rs/core/tests/suite/compact_remote_parity.rs index e221af2dc4..3c5f1606a5 100644 --- a/codex-rs/core/tests/suite/compact_remote_parity.rs +++ b/codex-rs/core/tests/suite/compact_remote_parity.rs @@ -7,7 +7,6 @@ use std::path::PathBuf; use anyhow::Result; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_protocol::config_types::ServiceTier; use codex_protocol::models::PermissionProfile; @@ -791,7 +790,7 @@ fn replacement_history_from_rollout(path: &Path) -> Result { .map(str::trim) .filter(|line| !line.is_empty()) { - let Ok(entry) = serde_json::from_str::(line) else { + let Ok(entry) = codex_rollout::parse_rollout_line(line) else { continue; }; if let RolloutItem::Compacted(compacted) = entry.item diff --git a/codex-rs/core/tests/suite/compact_resume_fork.rs b/codex-rs/core/tests/suite/compact_resume_fork.rs index bce268a819..522ef8c047 100644 --- a/codex-rs/core/tests/suite/compact_resume_fork.rs +++ b/codex-rs/core/tests/suite/compact_resume_fork.rs @@ -18,7 +18,6 @@ use codex_core::config::Config; use codex_core::spawn::CODEX_SANDBOX_NETWORK_DISABLED_ENV_VAR; use codex_history::CodexHarnessMetadata; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_protocol::config_types::CollaborationMode; use codex_protocol::config_types::ModeKind; use codex_protocol::config_types::Settings; @@ -94,7 +93,7 @@ fn seed_first_checkpoint_harness_metadata(path: &Path, retained_text: &str) -> R let mut lines = std::fs::read_to_string(path)? .lines() .filter(|line| !line.trim().is_empty()) - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()?; let replacement_history = lines .iter_mut() @@ -127,7 +126,7 @@ fn assert_latest_checkpoint_retains_harness_metadata( let replacement_history = std::fs::read_to_string(path)? .lines() .filter(|line| !line.trim().is_empty()) - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()? .into_iter() .rev() diff --git a/codex-rs/core/tests/suite/fork_thread.rs b/codex-rs/core/tests/suite/fork_thread.rs index cc0b3b24f2..ad9434a50a 100644 --- a/codex-rs/core/tests/suite/fork_thread.rs +++ b/codex-rs/core/tests/suite/fork_thread.rs @@ -7,7 +7,6 @@ use codex_core::parse_turn_item; use codex_history::InitialHistory; use codex_history::ResumedHistory; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_protocol::ThreadId; use codex_protocol::items::TurnItem; use codex_protocol::mcp::ClientMcpExtensions; @@ -311,7 +310,7 @@ fn read_rollout_items(path: &std::path::Path) -> Vec { let parse_json_message = format!("failed to parse rollout JSON line `{line}`"); let v: serde_json::Value = serde_json::from_str(line).expect(&parse_json_message); let parse_line_message = format!("failed to parse rollout line `{line}`"); - let rl: RolloutLine = serde_json::from_value(v).expect(&parse_line_message); + let rl = codex_rollout::decode_rollout_line(v).expect(&parse_line_message); match rl.item { RolloutItem::SessionMeta(_) => {} other => items.push(other), diff --git a/codex-rs/core/tests/suite/guardian_review.rs b/codex-rs/core/tests/suite/guardian_review.rs index 80695a8521..cf6f0642b2 100644 --- a/codex-rs/core/tests/suite/guardian_review.rs +++ b/codex-rs/core/tests/suite/guardian_review.rs @@ -22,7 +22,6 @@ use codex_extension_api::ToolStartInput; use codex_features::CurrentTimeSource; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_protocol::ThreadId; use codex_protocol::config_types::ApprovalsReviewer; @@ -618,7 +617,7 @@ async fn guardian_session_prewarms_and_is_reused_for_first_review( test.codex.shutdown_and_wait().await?; let guardian_rollout = fs::read_to_string(guardian_rollout_path)? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::>>()?; assert_eq!( guardian_rollout.iter().find_map(|line| match &line.item { diff --git a/codex-rs/core/tests/suite/hooks.rs b/codex-rs/core/tests/suite/hooks.rs index 119697581c..090ab49a0a 100644 --- a/codex-rs/core/tests/suite/hooks.rs +++ b/codex-rs/core/tests/suite/hooks.rs @@ -16,7 +16,6 @@ use codex_core::config::Constrained; use codex_core::config::ThreadStoreConfig; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_model_provider_info::ModelProviderInfo; use codex_model_provider_info::built_in_model_providers; use codex_plugin::PluginHookSource; @@ -1145,7 +1144,7 @@ fn rollout_hook_prompt_texts(text: &str) -> Result> { if trimmed.is_empty() { continue; } - let rollout: RolloutLine = serde_json::from_str(trimmed).context("parse rollout line")?; + let rollout = codex_rollout::parse_rollout_line(trimmed).context("parse rollout line")?; if let RolloutItem::ResponseItem(envelope) = rollout.item && let ResponseItem::Message { role, content, .. } = envelope.item && role == "user" diff --git a/codex-rs/core/tests/suite/image_rollout.rs b/codex-rs/core/tests/suite/image_rollout.rs index 76b7bc0b9e..3316742feb 100644 --- a/codex-rs/core/tests/suite/image_rollout.rs +++ b/codex-rs/core/tests/suite/image_rollout.rs @@ -4,7 +4,6 @@ use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use codex_core::TurnInputRequest; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_protocol::config_types::CollaborationMode; use codex_protocol::config_types::ModeKind; use codex_protocol::config_types::Settings; @@ -50,7 +49,7 @@ fn find_user_message_with_image(text: &str) -> Option { if trimmed.is_empty() { continue; } - let rollout: RolloutLine = match serde_json::from_str(trimmed) { + let rollout = match codex_rollout::parse_rollout_line(trimmed) { Ok(rollout) => rollout, Err(_) => continue, }; @@ -321,7 +320,7 @@ async fn resumed_history_only_emits_resize_notices_for_new_images() -> anyhow::R let mut rollout_lines = fs::read_to_string(&rollout_path)? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::>>()?; let historical_content = rollout_lines .iter_mut() @@ -501,7 +500,7 @@ async fn resumed_history_only_emits_resize_notices_for_new_images() -> anyhow::R .lines() .skip(existing_rollout_lines) .filter_map(|line| { - let RolloutItem::ResponseItem(envelope) = serde_json::from_str::(line) + let RolloutItem::ResponseItem(envelope) = codex_rollout::parse_rollout_line(line) .expect("new rollout line should deserialize") .item else { diff --git a/codex-rs/core/tests/suite/items.rs b/codex-rs/core/tests/suite/items.rs index e65d51059a..a3d3d44857 100644 --- a/codex-rs/core/tests/suite/items.rs +++ b/codex-rs/core/tests/suite/items.rs @@ -4,7 +4,6 @@ use anyhow::Ok; use codex_core::TurnInputRequest; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_protocol::config_types::CollaborationMode; use codex_protocol::config_types::ModeKind; use codex_protocol::config_types::Settings; @@ -361,7 +360,7 @@ async fn web_search_item_is_emitted() -> anyhow::Result<()> { let rollout = std::fs::read_to_string(rollout_path)?; let persisted_completion = rollout .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()? .into_iter() .find_map(|line| match line.item { diff --git a/codex-rs/core/tests/suite/mcp_tool_exposure.rs b/codex-rs/core/tests/suite/mcp_tool_exposure.rs index 2f6e6e8316..b554904652 100644 --- a/codex-rs/core/tests/suite/mcp_tool_exposure.rs +++ b/codex-rs/core/tests/suite/mcp_tool_exposure.rs @@ -13,7 +13,6 @@ use codex_extension_api::ThreadLifecycleContributor; use codex_extension_api::ThreadStartInput; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_mcp::CODEX_APPS_MCP_SERVER_NAME; use codex_mcp::McpResourceClient; use codex_protocol::capabilities::CapabilityRootLocation; @@ -1003,7 +1002,7 @@ async fn initially_empty_deferred_tool_world_state_is_not_rendered_or_persisted( let world_states = tokio::fs::read_to_string(rollout_path) .await? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::>>()? .into_iter() .filter_map(|line| match line.item { @@ -1046,7 +1045,7 @@ async fn deferred_tool_world_state_survives_resume_without_duplicate_updates() - let persisted_tools = tokio::fs::read_to_string(&rollout_path) .await? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::>>()? .into_iter() .filter_map(|line| match line.item { diff --git a/codex-rs/core/tests/suite/mcp_turn_metadata.rs b/codex-rs/core/tests/suite/mcp_turn_metadata.rs index ff322707c2..d3d5ecf07f 100644 --- a/codex-rs/core/tests/suite/mcp_turn_metadata.rs +++ b/codex-rs/core/tests/suite/mcp_turn_metadata.rs @@ -8,7 +8,6 @@ use codex_core::TurnInputRequest; use codex_core::config::Config; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_protocol::approvals::ElicitationRequest; use codex_protocol::config_types::ApprovalsReviewer; use codex_protocol::config_types::CollaborationMode; @@ -1229,7 +1228,7 @@ async fn apps_default_writes_prompts_for_writes_but_not_reads() -> Result<()> { let persisted_hints = tokio::fs::read_to_string(rollout_path) .await? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::>>()? .into_iter() .filter_map(|line| match line.item { diff --git a/codex-rs/core/tests/suite/model_switching.rs b/codex-rs/core/tests/suite/model_switching.rs index 6cfc82a658..0ed6cdf0e0 100644 --- a/codex-rs/core/tests/suite/model_switching.rs +++ b/codex-rs/core/tests/suite/model_switching.rs @@ -6,7 +6,6 @@ use codex_core::TurnInputRequest; use codex_core::config::Constrained; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_models_manager::bundled_models_response; use codex_models_manager::manager::RefreshStrategy; @@ -442,7 +441,7 @@ async fn model_change_appends_model_instructions_developer_message() -> Result<( let rollout_path = test.codex.rollout_path().expect("rollout path"); let model_states = std::fs::read_to_string(rollout_path)? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::>>()? .into_iter() .filter_map(|line| match line.item { diff --git a/codex-rs/core/tests/suite/pending_input.rs b/codex-rs/core/tests/suite/pending_input.rs index a636b2c248..c7ea6fe73b 100644 --- a/codex-rs/core/tests/suite/pending_input.rs +++ b/codex-rs/core/tests/suite/pending_input.rs @@ -12,7 +12,6 @@ use codex_extension_items::ExtensionItem; use codex_extension_items::sleep::SleepItem; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_protocol::AgentPath; use codex_protocol::config_types::CollaborationMode; @@ -721,7 +720,7 @@ async fn any_new_input_interrupts_sleep() { .expect("read rollout"); let persisted_sleep_items = rollout .lines() - .filter_map(|line| serde_json::from_str::(line).ok()) + .filter_map(|line| codex_rollout::parse_rollout_line(line).ok()) .filter_map(|line| match line.item { RolloutItem::EventMsg(EventMsg::ItemCompleted(event)) => match event.item { TurnItem::Extension(ExtensionItem::Sleep(item)) => Some(item), diff --git a/codex-rs/core/tests/suite/remote_env.rs b/codex-rs/core/tests/suite/remote_env.rs index 5484c90d12..532bad26f9 100644 --- a/codex-rs/core/tests/suite/remote_env.rs +++ b/codex-rs/core/tests/suite/remote_env.rs @@ -38,7 +38,6 @@ use codex_extension_api::WorldStateContributionInput; use codex_extension_api::WorldStateSectionContribution; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_http_client::HttpClientFactory; use codex_http_client::OutboundProxyPolicy; use codex_network_proxy::NetworkProxyConfig; @@ -1078,7 +1077,7 @@ async fn deferred_executor_promotes_primary_environment_when_startup_completes() let rollout = fs::read_to_string(test.codex.rollout_path().context("rollout path")?)?; let world_state_patch = rollout .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::>>()? .into_iter() .filter_map(|line| match line.item { @@ -2766,7 +2765,7 @@ async fn deferred_executor_compaction_preserves_then_updates_environment_once() let rollout = fs::read_to_string(rollout_path)?; let world_state_items = rollout .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::>>()? .into_iter() .filter_map(|line| match line.item { diff --git a/codex-rs/core/tests/suite/review.rs b/codex-rs/core/tests/suite/review.rs index 830fd1c512..5b48bb7700 100644 --- a/codex-rs/core/tests/suite/review.rs +++ b/codex-rs/core/tests/suite/review.rs @@ -7,7 +7,6 @@ use codex_core::find_thread_path_by_id_str; use codex_exec_server::CreateDirectoryOptions; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_protocol::config_types::ApprovalsReviewer; use codex_protocol::config_types::CollaborationMode; @@ -209,7 +208,7 @@ async fn review_op_emits_lifecycle_and_review_output() { .lines() .filter(|line| !line.trim().is_empty()) .find_map(|line| { - let rollout_line: RolloutLine = serde_json::from_str(line).expect("rollout line"); + let rollout_line = codex_rollout::parse_rollout_line(line).expect("rollout line"); match rollout_line.item { RolloutItem::SessionMeta(session_meta) => Some(session_meta.meta.id.to_string()), _ => None, @@ -251,7 +250,7 @@ async fn review_op_emits_lifecycle_and_review_output() { continue; } let v: serde_json::Value = serde_json::from_str(line).expect("jsonl line"); - let rl: RolloutLine = serde_json::from_value(v).expect("rollout line"); + let rl = codex_rollout::decode_rollout_line(v).expect("rollout line"); if let RolloutItem::ResponseItem(envelope) = rl.item && let ResponseItem::Message { role, content, .. } = envelope.item { @@ -763,8 +762,8 @@ async fn review_uses_updated_turn_permissions_and_approval_policy() { let review_session_cwd = review_rollout .lines() .find_map(|line| { - let rollout_line: RolloutLine = - serde_json::from_str(line).expect("review rollout line should be valid"); + let rollout_line = codex_rollout::parse_rollout_line(line) + .expect("review rollout line should be valid"); match rollout_line.item { RolloutItem::SessionMeta(session_meta) => Some(session_meta.meta.cwd), _ => None, @@ -775,8 +774,8 @@ async fn review_uses_updated_turn_permissions_and_approval_policy() { let review_context = review_rollout .lines() .filter_map(|line| { - let rollout_line: RolloutLine = - serde_json::from_str(line).expect("review rollout line should be valid"); + let rollout_line = codex_rollout::parse_rollout_line(line) + .expect("review rollout line should be valid"); match rollout_line.item { RolloutItem::TurnContext(turn_context) => Some(turn_context), _ => None, @@ -1229,7 +1228,7 @@ async fn review_input_isolated_from_parent_history() { continue; } let v: serde_json::Value = serde_json::from_str(line).expect("jsonl line"); - let rl: RolloutLine = serde_json::from_value(v).expect("rollout line"); + let rl = codex_rollout::decode_rollout_line(v).expect("rollout line"); if let RolloutItem::ResponseItem(envelope) = rl.item && let ResponseItem::Message { role, content, .. } = envelope.item && role == "user" diff --git a/codex-rs/core/tests/suite/settings_commits.rs b/codex-rs/core/tests/suite/settings_commits.rs index 91d328a144..f1ee6444ec 100644 --- a/codex-rs/core/tests/suite/settings_commits.rs +++ b/codex-rs/core/tests/suite/settings_commits.rs @@ -241,7 +241,7 @@ async fn compaction_checkpoints_settings_changed_during_its_model_request() -> R let rollout_path = test.session_configured.rollout_path.expect("rollout path"); let rollout: Vec = std::fs::read_to_string(rollout_path)? .lines() - .map(serde_json::from_str) + .map(codex_rollout::parse_rollout_line) .collect::>()?; let checkpoint = rollout .iter() diff --git a/codex-rs/core/tests/suite/spawn_agent_description.rs b/codex-rs/core/tests/suite/spawn_agent_description.rs index 317ce27af5..dab508f9db 100644 --- a/codex-rs/core/tests/suite/spawn_agent_description.rs +++ b/codex-rs/core/tests/suite/spawn_agent_description.rs @@ -6,7 +6,6 @@ use codex_core::config::AgentRoleConfig; use codex_core::config::Config; use codex_features::Feature; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_login::CodexAuth; use codex_models_manager::manager::RefreshStrategy; use codex_models_manager::manager::SharedModelsManager; @@ -454,7 +453,7 @@ async fn multi_agent_v2_cold_resume_refreshes_legacy_usage_hints_once( let mut removed_usage_hint_presence = false; let legacy_rollout = std::fs::read_to_string(&rollout_path)? .lines() - .map(serde_json::from_str::) + .map(codex_rollout::parse_rollout_line) .collect::, _>>()? .into_iter() .map(|mut line| { diff --git a/codex-rs/core/tests/suite/token_usage_rollout.rs b/codex-rs/core/tests/suite/token_usage_rollout.rs index db9e6c50e1..04def592e9 100644 --- a/codex-rs/core/tests/suite/token_usage_rollout.rs +++ b/codex-rs/core/tests/suite/token_usage_rollout.rs @@ -2,7 +2,6 @@ use anyhow::Result; use codex_history::RolloutItem; -use codex_history::RolloutLine; use codex_protocol::SessionId; use codex_protocol::protocol::TokenUsageRecord; use core_test_support::responses::ev_assistant_message; @@ -21,7 +20,7 @@ fn token_usage_records(path: &std::path::Path) -> Vec { std::fs::read_to_string(path) .expect("read rollout") .lines() - .filter_map(|line| serde_json::from_str::(line).ok()) + .filter_map(|line| codex_rollout::parse_rollout_line(line).ok()) .filter_map(|line| match line.item { RolloutItem::TokenUsageRecord(record) => Some(record), _ => None, diff --git a/codex-rs/exec/src/lib.rs b/codex-rs/exec/src/lib.rs index 68cde2368b..49ecf1edc9 100644 --- a/codex-rs/exec/src/lib.rs +++ b/codex-rs/exec/src/lib.rs @@ -1567,7 +1567,7 @@ async fn parse_latest_turn_context_cwd(path: &Path) -> Option { tokio::task::spawn_blocking(move || { let reader = codex_rollout::open_rollout_seekable_reader(&path).ok()?; let mut scanner = codex_rollout::ReverseJsonlScanner::new(reader).ok()?; - while let Some(outcome) = scanner.scan_next::().ok()? { + while let Some(outcome) = scanner.scan_next_rollout_line().ok()? { if let codex_rollout::ScanOutcome::Parsed(RolloutLine { item: RolloutItem::TurnContext(item), .. diff --git a/codex-rs/history/src/lib.rs b/codex-rs/history/src/lib.rs index c74fc5a2ab..74e5b01c0a 100644 --- a/codex-rs/history/src/lib.rs +++ b/codex-rs/history/src/lib.rs @@ -232,7 +232,11 @@ impl From for ResponseItem { } } -#[derive(Serialize, Deserialize, Clone, JsonSchema)] +/// One persisted rollout JSONL record. +/// +/// This intentionally does not implement Deserialize: JSONL readers must use +/// codex_rollout's canonical parser so nested decimal values survive the flattened envelope. +#[derive(Serialize, Clone, JsonSchema)] pub struct RolloutLine { pub timestamp: String, #[serde(default, skip_serializing_if = "Option::is_none")] diff --git a/codex-rs/history/src/tests.rs b/codex-rs/history/src/tests.rs index 46b1fd28d7..825244db87 100644 --- a/codex-rs/history/src/tests.rs +++ b/codex-rs/history/src/tests.rs @@ -53,7 +53,11 @@ fn response_item_rollout_line_preserves_shape() -> Result<()> { }, }); - let line = serde_json::from_value::(legacy_line.clone())?; + let line = RolloutLine { + timestamp: "2025-01-03T12:00:00.000Z".to_string(), + ordinal: Some(7), + item: serde_json::from_value(legacy_line.clone())?, + }; let RolloutItem::ResponseItem(envelope) = &line.item else { panic!("expected response item"); }; @@ -93,8 +97,8 @@ fn response_item_envelope_stores_metadata_beside_rollout_payload() -> Result<()> ); assert_eq!(serialized["payload"].get("metadata"), None); - let restored = serde_json::from_value::(serialized)?; - let RolloutItem::ResponseItem(envelope) = restored.item else { + let restored = serde_json::from_value(serialized)?; + let RolloutItem::ResponseItem(envelope) = restored else { panic!("expected response item"); }; assert_eq!( @@ -157,7 +161,7 @@ fn response_item_envelope_preserves_harness_authored_configuration_provenance() #[test] /// Keeps future metadata fields from making older binaries reject persisted items. fn response_item_envelope_ignores_unknown_harness_metadata_fields() -> Result<()> { - let line = serde_json::from_value::(json!({ + let line = serde_json::from_value(json!({ "timestamp": "2025-01-03T12:00:00.000Z", "ordinal": 7, "type": "response_item", @@ -174,7 +178,7 @@ fn response_item_envelope_ignores_unknown_harness_metadata_fields() -> Result<() }, }))?; - let RolloutItem::ResponseItem(envelope) = line.item else { + let RolloutItem::ResponseItem(envelope) = line else { panic!("expected response item"); }; assert_eq!(envelope.metadata, Some(CodexHarnessMetadata::default())); diff --git a/codex-rs/rollout/src/lib.rs b/codex-rs/rollout/src/lib.rs index 70abb519c7..71d67cc63a 100644 --- a/codex-rs/rollout/src/lib.rs +++ b/codex-rs/rollout/src/lib.rs @@ -46,7 +46,9 @@ pub(crate) use codex_protocol::protocol; /// Remove it once Serde supports format-specific buffering. pub fn decode_rollout_line(value: Value) -> serde_json::Result { let Value::Object(mut fields) = value else { - return serde_json::from_value(value); + return Err(serde_json::Error::custom( + "rollout line must be a JSON object", + )); }; let timestamp = fields .remove("timestamp") @@ -66,6 +68,16 @@ pub fn decode_rollout_line(value: Value) -> serde_json::Result { }) } +/// Parses a persisted JSONL rollout record through the canonical JSON decoder. +pub fn parse_rollout_line(line: &str) -> serde_json::Result { + serde_json::from_str::(line).and_then(decode_rollout_line) +} + +/// Parses persisted JSONL rollout record bytes through the canonical JSON decoder. +pub fn parse_rollout_line_bytes(bytes: &[u8]) -> serde_json::Result { + serde_json::from_slice::(bytes).and_then(decode_rollout_line) +} + pub const SESSIONS_SUBDIR: &str = "sessions"; pub const ARCHIVED_SESSIONS_SUBDIR: &str = "archived_sessions"; pub static INTERACTIVE_SESSION_SOURCES: LazyLock> = LazyLock::new(|| { diff --git a/codex-rs/rollout/src/list.rs b/codex-rs/rollout/src/list.rs index 8805d6433c..29dde5a78e 100644 --- a/codex-rs/rollout/src/list.rs +++ b/codex-rs/rollout/src/list.rs @@ -20,7 +20,6 @@ use super::SESSIONS_SUBDIR; use super::compression; use super::rollout_file_name::RolloutFileName; use crate::RolloutItem; -use crate::RolloutLine; use crate::protocol::EventMsg; use crate::state_db; use codex_file_search as file_search; @@ -1128,7 +1127,7 @@ async fn read_head_summary(path: &Path, head_limit: usize) -> io::Result = serde_json::from_str(trimmed); + let parsed = crate::parse_rollout_line(trimmed); let rollout_line = match parsed { Ok(rollout_line) => rollout_line, Err(_) => { @@ -1241,7 +1240,7 @@ pub async fn read_head_for_summary(path: &Path) -> io::Result(trimmed) { + if let Ok(rollout_line) = crate::parse_rollout_line(trimmed) { match rollout_line.item { RolloutItem::SessionMeta(session_meta_line) => { if let Ok(value) = serde_json::to_value(session_meta_line) { @@ -1300,7 +1299,7 @@ pub async fn read_session_meta_line(path: &Path) -> io::Result if trimmed.is_empty() { continue; } - let Ok(rollout_line) = serde_json::from_str::(trimmed) else { + let Ok(rollout_line) = crate::parse_rollout_line(trimmed) else { if let Ok(value) = serde_json::from_str::(trimmed) { crate::recorder::reject_unknown_thread_history_mode(&value)?; } diff --git a/codex-rs/rollout/src/ordinal.rs b/codex-rs/rollout/src/ordinal.rs index 3f002aec8a..29a151cfaf 100644 --- a/codex-rs/rollout/src/ordinal.rs +++ b/codex-rs/rollout/src/ordinal.rs @@ -10,7 +10,6 @@ use codex_protocol::protocol::HistoryPosition; use codex_protocol::protocol::ThreadHistoryMode; use crate::RolloutItem; -use crate::RolloutLine; use crate::reverse_jsonl_scanner::ReverseJsonlScanner; use crate::reverse_jsonl_scanner::ScanOutcome; @@ -67,7 +66,7 @@ pub(crate) fn ordinal_state_for_rollout( let mut scanner = ReverseJsonlScanner::new(file)?; let record = loop { - match scanner.scan_next::()? { + match scanner.scan_next_rollout_line()? { Some(ScanOutcome::Parsed(record)) => break record, Some(ScanOutcome::Rejected(_)) => continue, None => { @@ -111,7 +110,7 @@ fn read_history_metadata( if line.trim().is_empty() { continue; } - let record: RolloutLine = serde_json::from_str(line.as_str()).map_err(|error| { + let record = crate::parse_rollout_line(line.as_str()).map_err(|error| { io::Error::other(format!( "failed to parse first rollout record at {}: {error}", path.display() diff --git a/codex-rs/rollout/src/recorder_tests.rs b/codex-rs/rollout/src/recorder_tests.rs index e1714f814c..90c5523222 100644 --- a/codex-rs/rollout/src/recorder_tests.rs +++ b/codex-rs/rollout/src/recorder_tests.rs @@ -14,11 +14,14 @@ use codex_protocol::protocol::AgentMessageEvent; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::HistoryPosition; +use codex_protocol::protocol::RateLimitSnapshot; +use codex_protocol::protocol::RateLimitWindow; use codex_protocol::protocol::SandboxPolicy; use codex_protocol::protocol::SessionMeta; use codex_protocol::protocol::SessionMetaLine; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::ThreadHistoryMode; +use codex_protocol::protocol::TokenCountEvent; use codex_protocol::protocol::TurnContextItem; use codex_protocol::protocol::UserMessageEvent; use codex_protocol::security_risk::SecurityRiskScore; @@ -103,7 +106,7 @@ fn read_rollout_lines(path: &Path) -> std::io::Result> { fs::read_to_string(path)? .lines() .filter(|line| !line.trim().is_empty()) - .map(|line| serde_json::from_str(line).map_err(std::io::Error::other)) + .map(|line| crate::parse_rollout_line(line).map_err(std::io::Error::other)) .collect() } @@ -687,7 +690,7 @@ async fn recorder_materializes_on_flush_with_pending_items() -> std::io::Result< vec![Some(0), Some(1), Some(2)] ); let first_line = text.lines().next().expect("session metadata line"); - let session_meta: RolloutLine = serde_json::from_str(first_line)?; + let session_meta = crate::parse_rollout_line(first_line)?; let RolloutItem::SessionMeta(session_meta) = session_meta.item else { panic!("expected session metadata in rollout"); }; @@ -995,6 +998,68 @@ async fn resumed_paginated_rollout_continues_after_ordinal_gap() -> std::io::Res recorder.shutdown().await } +#[tokio::test] +async fn resumed_paginated_rollout_continues_after_decimal_token_count() -> std::io::Result<()> { + let home = TempDir::new().expect("temp dir"); + let config = test_config(home.path()); + let thread_id = ThreadId::new(); + let recorder = RolloutRecorder::new( + &config, + RolloutRecorderParams::new( + thread_id, + /*forked_from_id*/ None, + /*parent_thread_id*/ None, + SessionSource::Exec, + /*thread_source*/ None, + "test_originator".to_string(), + BaseInstructions::default(), + Vec::new(), + ) + .with_history_mode(ThreadHistoryMode::Paginated), + ) + .await?; + let rollout_path = recorder.rollout_path().to_path_buf(); + recorder + .record_canonical_items(&[RolloutItem::EventMsg(EventMsg::TokenCount( + TokenCountEvent { + info: None, + rate_limits: Some(RateLimitSnapshot { + limit_id: None, + limit_name: None, + primary: Some(RateLimitWindow { + used_percent: 0.0, + window_minutes: Some(60), + resets_at: Some(1_800_000_000), + }), + secondary: None, + credits: None, + individual_limit: None, + spend_control_reached: None, + plan_type: None, + rate_limit_reached_type: None, + normal_model_slug: None, + }), + }, + ))]) + .await?; + recorder.persist().await?; + recorder.shutdown().await?; + + let resumed = + RolloutRecorder::new(&config, RolloutRecorderParams::resume(rollout_path.clone())).await?; + resumed + .record_canonical_items(&[agent_message_item("after-resume")]) + .await?; + resumed.shutdown().await?; + + let ordinals = read_rollout_lines(&rollout_path)? + .into_iter() + .map(|line| line.ordinal) + .collect::>(); + assert_eq!(ordinals, vec![Some(0), Some(1), Some(2)]); + Ok(()) +} + #[tokio::test] async fn resumed_paginated_rollout_repairs_unsafe_tail() -> std::io::Result<()> { let valid_unterminated = serde_json::to_string(&RolloutLine { @@ -1034,7 +1099,7 @@ async fn resumed_paginated_rollout_repairs_unsafe_tail() -> std::io::Result<()> assert!(contents.ends_with('\n'), "{name} tail should be terminated"); let ordinals = contents .lines() - .filter_map(|line| serde_json::from_str::(line).ok()) + .filter_map(|line| crate::parse_rollout_line(line).ok()) .map(|line| line.ordinal) .collect::>(); assert_eq!( diff --git a/codex-rs/rollout/src/reverse_jsonl_scanner.rs b/codex-rs/rollout/src/reverse_jsonl_scanner.rs index 70919a7e61..9e0a836bee 100644 --- a/codex-rs/rollout/src/reverse_jsonl_scanner.rs +++ b/codex-rs/rollout/src/reverse_jsonl_scanner.rs @@ -4,6 +4,9 @@ use std::io::Seek; use std::io::SeekFrom; use serde::de::DeserializeOwned; +use serde_json::Value; + +use crate::RolloutLine; const READ_CHUNK_SIZE: usize = 64 * 1024; @@ -128,6 +131,17 @@ where } } + /// Scans the next rollout record through the canonical persisted JSON decoder. + pub fn scan_next_rollout_line(&mut self) -> io::Result>> { + Ok(self.scan_next::()?.map(|outcome| match outcome { + ScanOutcome::Parsed(value) => match crate::decode_rollout_line(value) { + Ok(line) => ScanOutcome::Parsed(line), + Err(error) => ScanOutcome::Rejected(error), + }, + ScanOutcome::Rejected(error) => ScanOutcome::Rejected(error), + })) + } + fn finish_record(&mut self) -> Option> where T: DeserializeOwned, diff --git a/codex-rs/rollout/src/search.rs b/codex-rs/rollout/src/search.rs index 152050f662..3b8f15860a 100644 --- a/codex-rs/rollout/src/search.rs +++ b/codex-rs/rollout/src/search.rs @@ -18,7 +18,6 @@ use super::SESSIONS_SUBDIR; use super::compression; use crate::ResponseItemEnvelope; use crate::RolloutItem; -use crate::RolloutLine; const MATCH_CONTEXT_BEFORE_CHARS: usize = 48; const MATCH_CONTEXT_AFTER_CHARS: usize = 96; @@ -248,7 +247,7 @@ fn case_insensitive_literal_regex(search_term: impl AsRef) -> io::Result Option { - let rollout_line = serde_json::from_str::(jsonl_line.trim()).ok()?; + let rollout_line = crate::parse_rollout_line(jsonl_line.trim()).ok()?; let text = conversation_text_from_item(&rollout_line.item)?; excerpt_around_match(text.as_str(), search_term) } diff --git a/codex-rs/rollout/src/tests.rs b/codex-rs/rollout/src/tests.rs index 915aa84f36..c0b21c4346 100644 --- a/codex-rs/rollout/src/tests.rs +++ b/codex-rs/rollout/src/tests.rs @@ -59,7 +59,7 @@ fn rollout_line_decoder_preserves_canonical_json_compatibility() -> Result<()> { for encoded in cases { let value = serde_json::from_str::(encoded)?; - let decoded = crate::decode_rollout_line(value.clone())?; + let decoded = crate::parse_rollout_line(encoded)?; let mut expected = value; if expected["type"] != "response_item" { expected diff --git a/codex-rs/thread-store/src/local/model_context.rs b/codex-rs/thread-store/src/local/model_context.rs index d27718e12f..6d88ad53cb 100644 --- a/codex-rs/thread-store/src/local/model_context.rs +++ b/codex-rs/thread-store/src/local/model_context.rs @@ -7,7 +7,6 @@ use codex_rollout::ModelContextScan; use codex_rollout::ModelContextScanProgress; use codex_rollout::ReverseJsonlScanner; use codex_rollout::RolloutItem; -use codex_rollout::RolloutLine; use codex_rollout::ScanOutcome; use super::LocalThreadStore; @@ -137,7 +136,7 @@ fn scan_model_context_from_lineage_blocking( Some(end_byte_offset) => ReverseJsonlScanner::new_at(file, end_byte_offset)?, None => ReverseJsonlScanner::new(file)?, }; - while let Some(outcome) = scanner.scan_next::()? { + while let Some(outcome) = scanner.scan_next_rollout_line()? { let ScanOutcome::Parsed(line) = outcome else { continue; }; diff --git a/codex-rs/thread-store/src/local/model_context_tests.rs b/codex-rs/thread-store/src/local/model_context_tests.rs index 23dfb7d0ba..6c32c8ea07 100644 --- a/codex-rs/thread-store/src/local/model_context_tests.rs +++ b/codex-rs/thread-store/src/local/model_context_tests.rs @@ -599,8 +599,8 @@ fn rollout_end_byte_offset(path: &Path, end_ordinal_exclusive: u64) -> u64 { let contents = std::fs::read(path).expect("read rollout"); let mut byte_offset = 0_u64; for line in contents.split_inclusive(|byte| *byte == b'\n') { - let parsed: RolloutLine = - serde_json::from_slice(line).expect("parse rollout line for byte offset"); + let parsed = codex_rollout::parse_rollout_line_bytes(line) + .expect("parse rollout line for byte offset"); if parsed.ordinal == Some(end_ordinal_exclusive) { return byte_offset; } diff --git a/codex-rs/thread-store/src/local/rollout_lineage_tests.rs b/codex-rs/thread-store/src/local/rollout_lineage_tests.rs index 4f4f762aec..6b8c88d118 100644 --- a/codex-rs/thread-store/src/local/rollout_lineage_tests.rs +++ b/codex-rs/thread-store/src/local/rollout_lineage_tests.rs @@ -327,7 +327,7 @@ fn rollout_end_byte_offset(path: &Path, end_ordinal_exclusive: u64) -> u64 { let end_byte_offset = bytes .split_inclusive(|byte| *byte == b'\n') .take_while(|line| { - serde_json::from_slice::(line) + codex_rollout::parse_rollout_line_bytes(line) .expect("parse rollout fixture") .ordinal .expect("paginated rollout ordinal") diff --git a/codex-rs/thread-store/src/local/rollout_migration/line_parser.rs b/codex-rs/thread-store/src/local/rollout_migration/line_parser.rs index 4d004a533c..b8a50d044d 100644 --- a/codex-rs/thread-store/src/local/rollout_migration/line_parser.rs +++ b/codex-rs/thread-store/src/local/rollout_migration/line_parser.rs @@ -39,7 +39,7 @@ pub(super) fn parse_legacy_rollout_value(mut value: Value) -> Result(&head_bytes).map_err(migration_error)?; + let mut head = codex_rollout::parse_rollout_line_bytes(&head_bytes).map_err(migration_error)?; let RolloutItem::SessionMeta(session_meta) = &mut head.item else { return Err(migration_error( "staged rollout head is not session metadata", diff --git a/codex-rs/thread-store/src/local/rollout_migration_tests.rs b/codex-rs/thread-store/src/local/rollout_migration_tests.rs index 7ed8810601..6ca78276ea 100644 --- a/codex-rs/thread-store/src/local/rollout_migration_tests.rs +++ b/codex-rs/thread-store/src/local/rollout_migration_tests.rs @@ -259,7 +259,7 @@ fn read_rollout(path: &Path) -> Vec { fs::read_to_string(path) .expect("read migrated rollout") .lines() - .map(|line| serde_json::from_str(line).expect("parse migrated rollout")) + .map(|line| codex_rollout::parse_rollout_line(line).expect("parse migrated rollout")) .collect() } diff --git a/codex-rs/thread-store/src/local/thread_history/read_tests.rs b/codex-rs/thread-store/src/local/thread_history/read_tests.rs index 798151489e..06eaaca703 100644 --- a/codex-rs/thread-store/src/local/thread_history/read_tests.rs +++ b/codex-rs/thread-store/src/local/thread_history/read_tests.rs @@ -1392,7 +1392,7 @@ fn rollout_end_byte_offset(path: &std::path::Path, end_ordinal_exclusive: u64) - let end_byte_offset = bytes .split_inclusive(|byte| *byte == b'\n') .take_while(|line| { - serde_json::from_slice::(line) + codex_rollout::parse_rollout_line_bytes(line) .expect("parse rollout fixture") .ordinal .expect("paginated rollout ordinal") diff --git a/codex-rs/thread-store/src/local/thread_history_materialization_tests.rs b/codex-rs/thread-store/src/local/thread_history_materialization_tests.rs index 3ed9c360d9..f0b71bbbb0 100644 --- a/codex-rs/thread-store/src/local/thread_history_materialization_tests.rs +++ b/codex-rs/thread-store/src/local/thread_history_materialization_tests.rs @@ -504,7 +504,7 @@ async fn paginated_realtime_items_materialize_separately_in_rollout_order() { let expected_rows = fs::read_to_string(rollout_path.as_path()) .expect("read canonical rollout") .lines() - .map(|line| serde_json::from_str::(line).expect("parse rollout line")) + .map(|line| codex_rollout::parse_rollout_line(line).expect("parse rollout line")) .filter_map(|line| match line.item { RolloutItem::RealtimeItem(item) => Some(( item.id, @@ -2675,7 +2675,7 @@ fn rollout_line_byte_offsets(path: &std::path::Path, ordinal: u64) -> (i64, i64) let mut start_byte_offset = 0; for line in bytes.split_inclusive(|byte| *byte == b'\n') { let end_byte_offset = start_byte_offset + line.len(); - if serde_json::from_slice::(line) + if codex_rollout::parse_rollout_line_bytes(line) .ok() .and_then(|line| line.ordinal) == Some(ordinal) diff --git a/codex-rs/tui/src/resume_picker_transcript_preview.rs b/codex-rs/tui/src/resume_picker_transcript_preview.rs index 8ef505c0a5..fb642f5362 100644 --- a/codex-rs/tui/src/resume_picker_transcript_preview.rs +++ b/codex-rs/tui/src/resume_picker_transcript_preview.rs @@ -22,7 +22,6 @@ use codex_protocol::ThreadId; use codex_protocol::protocol::EventMsg; use codex_rollout::ReverseJsonlScanner; use codex_rollout::RolloutItem; -use codex_rollout::RolloutLine; use codex_rollout::ScanOutcome; const MAX_TRANSCRIPT_PREVIEW_LINES: usize = 6; @@ -162,7 +161,7 @@ fn scan_legacy_transcript_preview( )? .with_max_record_bytes(MAX_LEGACY_TRANSCRIPT_PREVIEW_SCAN_BYTES); loop { - let outcome = match scanner.scan_next::() { + let outcome = match scanner.scan_next_rollout_line() { Ok(Some(outcome)) => outcome, Ok(None) => break, Err(error) if error.kind() == std::io::ErrorKind::UnexpectedEof => return Ok(None), diff --git a/codex-rs/tui/src/resume_picker_transcript_preview_tests.rs b/codex-rs/tui/src/resume_picker_transcript_preview_tests.rs index 61b99ed295..0c406185e9 100644 --- a/codex-rs/tui/src/resume_picker_transcript_preview_tests.rs +++ b/codex-rs/tui/src/resume_picker_transcript_preview_tests.rs @@ -18,6 +18,7 @@ use codex_protocol::protocol::AgentMessageEvent; use codex_protocol::protocol::ThreadRolledBackEvent; use codex_protocol::protocol::UserMessageEvent; use codex_rollout::CompactedItem; +use codex_rollout::RolloutLine; use core_test_support::responses; use pretty_assertions::assert_eq; use tempfile::tempdir;