From 64c9cde458b6b1c1d332577425b927f0436ab40b Mon Sep 17 00:00:00 2001 From: guinness-oai Date: Mon, 31 Aug 2026 22:04:14 +0000 Subject: [PATCH] Record realtime conversation history in Core (#41924) ## Why Realtime history should be recorded consistently for every Core host, including when no app-server event listener is attached. ## What changed - Move transcript segmentation, session boundaries, and backing-agent artifact promotion into Core for paginated threads. - Persist canonical realtime items through the thread store in event order and emit dedicated history lifecycle events for hosts to present. - Translate those Core events into the existing app-server realtime item notifications without app-server persisting the items a second time. ## Testing - Cover Core-only persistence across repeated sessions, ephemeral sessions, accepted and rejected steering, typed input ordering, and artifact promotion. - Verify app-server notifications correspond to the persisted timeline. GitOrigin-RevId: 7cbef14129d77f6d6d7099b733be91f5279c55f5 --- .../app-server-exports-stable.json.zst | Bin 144262 -> 144369 bytes codex-rs/app-server/README.md | 2 + .../app-server/src/bespoke_event_handling.rs | 36 ++ codex-rs/app-server/src/lib.rs | 2 - .../app-server/src/realtime_event_handling.rs | 91 ----- .../request_processors/thread_lifecycle.rs | 54 +-- .../src/request_processors/turn_processor.rs | 50 +-- codex-rs/app-server/src/thread_state.rs | 7 - .../tests/suite/v2/realtime_conversation.rs | 105 +++++- .../endpoint/realtime_websocket/methods.rs | 5 +- codex-rs/core/src/lib.rs | 1 + codex-rs/core/src/realtime_conversation.rs | 5 +- .../src/realtime_history.rs | 198 ++++------- .../core/src/realtime_history/presentation.rs | 129 +++++++ .../src/realtime_history_tests.rs | 317 ++++++++++++------ codex-rs/core/src/session/mod.rs | 32 ++ codex-rs/core/src/session/realtime_history.rs | 56 ++++ codex-rs/core/src/session/session.rs | 4 + codex-rs/core/src/session/tests.rs | 2 + codex-rs/core/src/session/turn_input_tests.rs | 60 ++++ .../core/tests/suite/realtime_conversation.rs | 278 +++++++++++++-- codex-rs/protocol/src/protocol.rs | 7 + 22 files changed, 970 insertions(+), 471 deletions(-) delete mode 100644 codex-rs/app-server/src/realtime_event_handling.rs rename codex-rs/{app-server => core}/src/realtime_history.rs (73%) create mode 100644 codex-rs/core/src/realtime_history/presentation.rs rename codex-rs/{app-server => core}/src/realtime_history_tests.rs (65%) create mode 100644 codex-rs/core/src/session/realtime_history.rs 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 c43dd8756d20160edb18af12999f99a1527b41a0..3ea880b4b6914e72af48494ed64ec515f435f3c2 100644 GIT binary patch delta 6188 zcmV+{7}Mv5=LqrV2!MnEv;yrLe^jpk=3!$a5OdoIPSZ5XEC4@DY)x7=Xqh@)Q^q0# z=>Ds9!wLWZ02~Ka09yc70R8{U7)AN8@(Vh7ry9L#)Hx{>7|OeY_xImAsC()!e@8JE4D3@wyDu(7 zhASxaz4cJ*;mibva;Sq1*G-aSm;*x~5J;e4;W2?g77xXOVIdjD^8+Jb0O4R*3I>dz zATS67h5!PA;35cuAP9n>2!bF8jW8qxRHT{I8lxcI^IQO!*}=v>T)4AQFbrd>t^}WK z^-;Y6E{VvkKg`?kf2SAveh|N4okdF<-F)E7_yOsOKWZ{%^ltSWe+ob=4SlQOoQ#+$ zKcyR$7MndjYA1Y0re6*{%(7cZU#uQH`*F2k%9Tr~p$6?Xm^$i=%GV1va89K>=(pK^ zJuz=CEG;mvl_}=8?P7SEznWBMGA%kO4JX~JsQHFNS23kLe|5+S<;-YY)8*iuvE{4L zleba14`GZi53ynwzwqpe#te*OY*lR&TAci!O!1U9*WaC0e}tgzM~Z6lf2o(>L-M2? zvHjp<=tHpqVS7%b@G!AP^V3v#SW)UhVS_6?d`HY}LG*P3 zBI;{wK=AMt*zkb-e`?Gn^?;w3u-c;G_fWCWRv>rnf08I(2BMVM%0ucMZW*@(6vI58 zYqO%h9`6eFZQfoWp!FFV8=Kn*%S?8`quz%u=F(X_IBfs# zy9`d=&959ER>0-6H4_76}C7SC; zb&`uie-&9K!wXR^I?6;Y`7pj{qLVumlDP@bgV`ktGse?kE;q}^d}7T_XXKjGJpPeU zcTEy@?3eT=TqhijgyjGG_;KA@i-QX^P~D}A{5aJ|VIGc2$w3|rZB$=Ub%Q(B>y=FL za=JZ-(w@U*Dbw>&yTvgyHbAw>%MK zW7`6sEuGfo?S}XmQ;G#$0UUa+&b>G8zz%_<`E7w)EMkG^zZNeJzMoeg1tKc|Gqo*t zkD-|u40>CTT|?rpto)~)s78=`8nG0K-+ppS)k7^mGkR6%LVh_RI{gFS&VqY@#w~6%olpSwk%1N926H9ZEJa{d zi&E(UR{D-8l_dez8~y<+%%nXNFpqdLf8!$-Er;y~R^7RR!KW#fGf0}vPD(qRt1P!z z^hRUY2W~FI%`Mh&C_*fxC?N)%xLG@RUC0f7QZx!ts*aI#QbY@6a7qjxL{f6_2U6`( zYGXJtfq*nDn7C@99!k>>teRNpSw|vq@sZ8Fhq1Jov4REVSjp+R7;eJ?)(ZZDe|U2K zu86n09WSC0_z(@*W(6Po(!_0*EGBOSqr=gGw=rHlJ=k7;6k0))P>$)}=O3*&{P6R| z$bT^TF%zM{nzAtJMijsH0VB9GCvtBP*Z0gg@8Wb27}#iV!UqypF8GOtkx*s1V4z=Q z7HChgCeFa!g2%{q3r_rWsqaYIf3=1*Z^`>>4e;L4sl)+0n11Ze&1R}$e>2~lXj@*M zgS#n(VmaJf6Oh)Dp-##bVvC_`}KjZ%|aj z&M9QdHvTK--_SA>u&(yq8@ze)tgO^!looq}Ly$NhI=4Ebtilpx+-gL-D24M1?RGiq zPPi`S=%se?G$24RO3bG35oYPAEsEq%I*5rfQnMI4>c9tvDQ%Yz7`00jci23dX;W7B zjR_Y~13m(!0NnR+&A=6le=?fjA@-9J^P!ktbtY!bUZ!87>kO+R+ain8V1jT!Bf=tN zAFoZkIt*WBgsEtVEbwxGDY@nEKTdnY*twn$)0Qyj47s-Z@CmYv)vx3A_5 zm7iq609BtKtRkRFsb#S7-QuFK$%eHI{x0{#*Bz|02G~D<-W$7ve?bue)`s@%D?SZF zqKtw=?MiMMRB*WqqdiWAJAnl<+{N1~zQpGf%)Z;93MI8;p;SSqxpn&-Op_ zA7JmA;!vuslfqYRR^A_Md@;371+9*0S0*T>@D0E!R0ulje*u7=Ai}KA4k2q<5NnHJ z9?gIs=6X2)e-_A3@Ps$nWLtyKYwmgFiR9^O5lUyd+u$Sba>=vTQHj~dF7z7G2ZKiB zTU|T;E!wdq4*E{WVxZrr6m{~rjIJSE?ZyT zI(~b~0(ugi`3Q$m`^7=nni9Z4 zDJ7zL|1VR*V**OQz1y@gaC~|Rs_j1VJqT#K96$oBV34b3suqiF+DUVMrp(D9Lqm~b3@I#4DmJP=0goXOniDGWGs*h?=yqi zozXAX(EzIt?OuSEG<^XItrO|pF_2-kFLAg>0Hm${8`BAJ)9C8mX`NiXON=*?o#?G^ z^)BWUwyP4?dyOxB5!-}Il;r={lR6X9%>ZzNe-9J5xFl37mq0!#U%H+MGnuD_>N5q# z&}W)lZ>tCdSeHVg!r_CXGLQXL#RQ^sy6 z0$?OMMtNCJVt}O$mY=Ffi>Uw(=7gjwl*s)GzM?wtXLfIwL!SdWJ02a)uud;kD&b*` ze`g}0KCJYHq>r?Tn7YHpM*!>lUJ_b4%Lb%&R7X$2*>yMqvg&H9h!dD}V&?o%PYsr_ zJeT}Mkfn-C-g&aPc?grC{~X@t8zgCl&sV(-gXcH2N=mh?VQtGCKgN2Abbn*EV;q3z z$hF@kk~k~AivG|4m`U{wxFq_VvfESYRZx!?vz^LLh->V81elX6 zS(jJ}j*0q;(!HOws{KbyOwm?fJ@IRPaveXj3Yh=T8-FO;3}E|QA}`M#hcKxt-4m{c zD(OUB7~m=uS65Q``DyO+7aTf=+`mT&gm)E&*Tn4#^2Sy5?YkEz| zk8={wr^oP8o(n!%NTj`h#f3g2f3odnfFPixry{3@K$_&GsvlQ1qS)i<@$ zXM^M9(-jRWJxwkB0drFa)!6Vov;zO_J}*n+n3;Q&a^AMq5B|;#T?qCvGr^mU^-^R^qX^>wONS7i` z2k@CHdi&mCFSEyWrpC{Gb<_>NYk`?r-SsW%PyW*x6zF2HndHE~E!sd=s-~D9VW+5I zJL;BJ$tMA7bjNh`pfak`IqzvC&)OG>h1tRRu6KMH8 zT|kW)ZK?C))*R%of8v)RK2$b&=ddz8u5&6`_oLG?7lcHOH%{ZwRPmDWyZQPw>#_uH zm7`YAQ*az9#K);_q#&x~^J_q78R(3@l#GKn_KMX87>BX?K?v6MmW-wi=hqwSi#UL+ z@IHRw`(tYoJm%*JikWhua^d52A$mGBH^H6ln&X!Z#GMb9e;^N38HmL(4n_F_khh1p zKMiY)J?;URMRY1|C=~uP1i8xRz0}g0(W5lD<8Fou*#-o7?5bP8lb}aZ_WfF!lJt8tr2!F$!eXe=a*XBFauH*_zq%88_3@ocBj>F~U~zx2*TI5uJLRyd(I?s3Stu0ISDS z2V*ncl7Q1#)O)ANM@D3oitQ30!Kl^sw974;IiJ}x59zL%5-GAKoZnw@P59BlQ1ghy? zb64?Z`re@8mkrWEaa%_S#ZX}{l~q{`jI_d_qDkxQ(8>hblj*Z;4*pt0zS~OPknVyr z5yt64uE=8%a=2Lq(xP@ubZ#APH_<6Y3+Xice+LPJ=|18r#Iw0yXXu?LPM z6%exQt7{7F@$U{8Zonv5z(eq604eVTO$1mVme-s*e2Ye~H9}-+t40tx;wDXi&D<7c zf2BIf?;6lO*d$>v%)lN$gy>>70`M3v*bv=zHp8o_ zRbXmOZL%`7TW{9X2K!)Bn?N6%+K#VgQ)?5MO>HDe+SHbQ*rryd^ft9Qu)+OvUH;zq z01B{VH@#uOa%(w|i$CXdQ!D>G#%Kb-6Nc6pi9TQx-G-pIm7b{`Zy!W>l?1mohg z`M-={d|A&ox>nvBmbR6qKYkQ(zXLR+MWXKcvaPt29d*Kp^+=SEfkJ6UHxSNtTTE-g zUc$g<{y_$2Vj>OjcMDMq^#@bXf2YjP59Kn{g9lBg$SQZiNAsS`H_9ye+vldBbHLd zgI+-IG77ph9yvkVw7y2eP#kgCOatNMx{XH@YpgU6U&Nvcy6f1S9sa7@>oX}SuP1lA zR5$Caqlcu3EYY7!Hf&u)?J0Sa&s4Wt9^9VE24Y+V|%FM1&e@FI>3=sENLA?=MLkE7WBaq2c9Lb8^E3P!PSvKa3q2x97Fv z7!h90D9ozeX5UW{fA;}T)=H>U$#g=Xg<2yk zFBR6O^xNV7cD?DqS9f*Yt?x3cDnB=-INy!T!xlM}iTmQq8jL%6t-=uGQu~k|NQnBs zQJG@NHIdp`PD-VhaWMiLe-I_5#|E>M)L9+8HE2q5tAG2Yn%yOO=VUBRS4HF92;`w1 zFpx$1OB|Sqy*slVjCN9OKG1Qz{g^F8g055;3%NgjCa349w>eGjcn46Ggi#2F9w737 zVJRP+4Qy+fuHs$ghaOX?NasNTELuD`gY+8hW*prH-DC{KadS%^e+IyH?@P_QaBxRw z-Xd*f1O)7HoqE<*SZ%fLd#dwKn_^U!bbA!*_aC$T`nZ+H*+x#St#Hi_B_hZ`Efy`ZX=!DR3fz7J z+&)+tYL7TZzy#eYHU$d5m@(Cc47gB*Rcg*@1uKmxHao7YR9!2Z!shvC8h||cw*Vz^ zgXrcCdOv4n!MGKt KQB4iM7b`6n!_{yA delta 6080 zcmV;x7eDCn=Lm-92!MnEv;yrLe*~-m2H{{M0F&F`$4C-MT7n>UHgziUCEL@t6x1T) zer+#1f{>1{661HlFija`#VlS!Mgw3g}2hcYvx@K4wuDc zUd0O7OU3%p&EWOB(MjIfe_lzJx{}<(`liOeNa)xO9UUe*#H4G|^CaZDkte{f>I}a0 zkf45xQCAbjHs3;*yzC$W0pZzo;8Fmhl2z7X1j3k76qkIH6#4gchkc5$I84pCG8>F# ztc{;Is+qWSO~$E4uO#Z6m5B@Gjp6QJQfheVsT*==J^35U=ZP8SPF)Vpdc^^1cm?tfgmFY zf*=Tjpa_B>2#qi#1XQA_4QuNHis#V)Ftg(uOPE(@Q^9P0Q2h!Z+bT%wzOl%FZXIN3 z`W|1f-UIZ)a!o8bf62{*Q2kDdCyuDel0mQ4(kLg#dL;B#!*enkpj;FRG@3^137ws= zK1{!E^iZ5`LHta586DuO;ZII1p(Y!&!+`2I$4IC13=QD7l^IbTs+ zo@y7fPjk9}H>PRB0>&ibUPjHA+_*~j_}QG)WnY?0*v|gbgNH-^Vdh?`~Zj#xDQPP?7`^N>FrxI?sRilZwal z7oQl9{rZaV>@QloDG$@&oQ=R~E|(4I-600^V|f=?3-7Q7#h=7qKSBdvsmx~=@u5isuX5@C!D*njSU6WY5%$D>j zual*Zgk<7}tf1?WJt4Gg-cwVeL*mXkZ#f$=Imk^!NH}Hb#9dS771UecCP%1NN%J0D#40R%w!z5^&gRs0Bne8z4(T zDz;PuuJ+S)KyP?eRv|5$_&q&MFg3yPK;u>rO>maJe?&*kf5TgVH66)OkD~PZjmqqZ zg0}NU-SEdzDV?;(304KG8CEc9x%W9hf7Mkb7|M;bFoXHu?AW!#CCaIac}^QfHmB00 zr>d9|hyK7gihyA4M8^@q3m}x*xWp*)i)uwmC(TwM^TSnWreFdH9-(mlJNLhu5eVqw zTqM)3G_7k~f{^!vf;AGY)9l#Qdl(bU`V%22{z=O-O>*+seb0(QaF5_5243__)PNCG4;ZW!f5kBCOgj$gmpYG@A66067Ab3D@IlVf3U=CMDm6v*XM3H= zpFqe`4&?wtmE)R%tP#CMk7dO@Sl|k5ai{u9sDH*)qe^-HPHwu}GjO(_58(txReKx^ z7WhdSSuo9Aovv8Zh9e86m#c_pHJG?YZ#t8@H?0QK1>%y8>gFw|7exZ^O4ksn zLNc&W0+q)E{LK5ZROl41MGFC6HArUFdfnvq^mpjaFxBb}-xj*p(y^JT0o{^DP%0X^ z@L_AoZ$n2v;o(f{NrH000RkTZxm{CUlk@Kw(Pj92!j$HM+nQUd)@VO&1&1~xNynj0 z{`m;cr9(f0KpCFde_7vI;I6DeX28*b?(8K(@nO#mkj7R0Bdx5fRlJ(l^8RaYh(O11l> zkXB7SUK&y!sf5ecT|^T(-n5Y1Z&5>uBZ2B+B|C zkb5W9e5YNZ4H4e3@P2Wey>cknUzZH~nF55Bf-#yH{v`-T>`&>D zcRd^ft4{CkYPUpq2DxT;A0Pp3@(W+k(@tu4{%^%@9 zcr8TDfA&s2?b)z9=kl-N+>kdeL$IdYz^ld;6CX&sizJikVl#8-uHwrT;v02`HvK=0 zH0%Fa*7@fy5;&x)lX&|h0-&nz`?-m0G;(ye3Qru}{iieloCvJn=q_ikHI5RynZqow zh;0rR#K`5=QwS59$^hwu#|)UbB-AR7K&&SRf4ZLN8Cj`>S|_c?uugipx>HdIkT11U z&pN4AoumkZc@=X~oyuYBjllyQB%!h)f^5_kiG!T|2`RYS3^BW1>Spa~Y@J^&f&or2 z_%sz?7CHeRY!#AvC>^gQRB&|=Pl%pBhgNrInm#%nV*gUA1;Uez=diMH=<5yUkF=ec zf10a}{r=U~q>QheB|cKeiHaZ9+4T{^V+@rHfEti|V&(k7PmTJ<+!OfK8bOLp%Q>Ok zY)d)B{{wdPO;Q5$?5kRh9=8rHXr+qTu(g~m14x-cTHi?2@!P*KsjTKiOU`#*#Z_V( z%u^{JT$X;0OB;H-DvD9R)HVG^BIdRte{O$;B}iowsEnC10<_633D3}kMl011_m6rU>I5_pRy3+20kiV3`4YcBVTHpn;#)_%1a;3)^ydwdy0uqdNMah}m z<|E&uzEWjfhhPzuGF-|ipC@O);cCNpCSrE(fDM&x769A3Tr6W^IVa84@$WvdDZuLo z36>Oqz4%z)pmj552n+BPMs3uqe{se~sWxRb$#?il{v1?H|3Co@Ycz7GoUZ2hKbdZ% z1|)v19gI^2Z>HD%23r?8MNP@rY0YqyUtvWXK8PT)*OT?t+l_WJ*>ISA@}dDxFGDM4 zz;e{A+UQ-d)`Wj^$rDG!n>hmGL^#^Kdb?ZZ)W3uJko!GBl>s2h+O)WFe|IGO#lP=; zWrV-cJ!2M$Ht{LTM>|FCyR!8MNEBU&Tr$6lI?o{9Uj}IsXXv<# z95A)YzO9{-IN6!O_D47crjMKNLzuc2XKLzB;moojyA|PCR_X^!5g)yomaLhjU-QRU z$!I&6WhMxq-7#lJIIz8ve?v?dX#3J8(~9%l-{ull1ARPiAblgI15hV*;M>C>{dR)K zNu8S8x~YkLhlmFls_QQ-n&eG5VcMp!&nxxY!j>ACVK`cIt)(gk)vkAU4*QC_ph1`; z7K5NB2qm+if9rmw?aFHn0Q_aINYNyQ3J!5-frvT$!~XX78Ie&-e~(xBy$q)?s(Di0 zH}xo_)6~7ngA^M0u{ldDHqx1J(nIE$T7s;;a9h$d(9t|+?>5$47ub1|CBw3Fz~j0JY=6TyJM3@)K(RFld?OWoHhz^7gms$=lWoJa>W8>Z%Ydv5B)2@ zqhE$Jp(4_|*_F{le;olzZ;x?;@L!L3g$KBQAJjPThQnXQd&I9}{-BjSnp0JD$}rEX zn0z@On5-x9!6ywq0-fDKXGmo-9EWL9Iz2FnH{N%09<9sqXSSRdX?(=De^d>(6e{uNO>A$x}xW5?LN@BX7 zJ`fT!heBr#6*$w`?w<~nf&MS7Iv(KcqmKvEf|0<2xd*-Ffe>~yRP3f$yf}-O$N}KOQjH_06vDtr7Q+NiCXo)?l z#UmIAO218JfA}NWtER?XsN!HtF4!{fgjrA8h$=dM)Dh})ILG6t0gaBU4r4Rjk^w#_ z3V0i@FG}El_%{Wxok$UvV3`^E2oam3a3#UBY)aPtuvN2a1)bsQIjhe{7!`#Q1?5Y2*O3_2(gpg+|HyD~~*KEINO9C)cE;0YHYLoH(8`5V4} zQ~~%5e^uu=QebriKE!R&o(h_!uB=`*b9>z6DL%@}+CdCG-i|et4-CycSSo&bp>?C2 zUnE5b1)Ufn6!(O^yj#tpcBByFjV80RV=EliO*3CJQ)pl9{%+%PgSrdOvKC(#^57kL z$!N{gk%p+Fapk_fatKe7_35xr3EMx)gM@xoJ^>h?_5(F z#&LHr-&{^W1lD_(8UD~E*cD(6SjMEXPgAsYqN!$TTY^_{v@E;cG`EFX<(lMIBvy@U zf6@%bbFjx>LP&Ax04Se=H^kBb2<#Q8Ro-=h+VT9!XA9JlO<|yxC&L-2&0*R=ZHId6 zTyG83R`z+ImW%!awRNgOP}?99g4#1l5!9wXNKnf&(-PFG(AN6nP#za0fD;hxFZ7^g z@bvwutp7X|)H43`zm45d_($o~M+k?Qf0uvsuPwn12CtrWw|T&c)BAM|w8^pC<6LEy z1py8JpS3CurODhWXs$iyCU?GG>znBWKgYxNM4ByH*WwZdo}4QugnK;VFu3CyoGQw7 z8n?!Aik#?%s`3O%XD$`s?twE5c1l8?*V@VknP?{r&%;PU1~p1!rRh2gZ!zH%f7}wb zYDNh%u!SPhfMK!_uuy+Y1*OPxe#kImA3XT_6zPAj_6X`-_(thVz!w?HaUM;HAZ|Yr zq!lUiy`d!&hku39hO6?*d3ka8pN4OB2xMSm|FV(_jf20Zg%kUPaA^^xYcTd6Pxhc_&rv3sU{qs_4*l;IM zx$J{3lSg1_n|5)~Fytc+%NqzM_x*U6A1#&U$%`nfpgW5Vsl#9Ob$uoX<@My&m+GdM zb@W&ik%j+r$pEbjiajM>@|o%u%j2i#umR8Wf>fUnqV!(?ncC?=d`xZKe-<96|D*k-j?9 ztLGow`{9+tE}~oU!(nvD8i(2VwGclRWREA^mQ(~?dFcT0N)@(|OXj=sRRrzxFA&QP z>LOPp4ahPb)QzGPSrKm>f7A^Yg@d|+_f{*HKMLCo>e9C6SmP7bar1Sx`=t4W&%Lx2 zCT-u7-xEB(ZwtF=<>opbMmgz@dl2zwxl|;M8R%e&M7RCu>kvsweA&c=CrcXIT&4Qj)1nIjU7~Ax&2%oE$nPrJvVXf~=sl ze+&J}sqQ?$N}nzXWbXq%x8II>Q~3fSY3hyq=w=r1Ld;F5Lg(|LRkJm#rw~wN<)!%3 zr<^+Kzg=(o{?%PwFWbwz%<9CllPR3Q;~RCo)?!(fg8-8(jh|zLCYlwlOSg z_m^11OuXBfm$1-D_3{wS@!siFVgwzc3==+nyh_fFqYj}OzvB_0DvZB|Ej=IP0ZS<# zI1TLfveb(A*oVP0kI0Ns0xGn4B!hSx?VmW_1|3XH`*8>l?0@CWb*q}1H|XFd&vZoE zGZ97DQ@Y6=p(wW1Q17YGKb?sgS<>NQ`0cW}+c_iL^{op|BZlngq)HL*YaK$^swfWt z1~pcX7fu51FHYpYGX*ujvBK29XbKXf#8rGM?}Hbqz{)ZA9HBdu#|!~s|?l&#Z+2vs9X1NA3b6w=ZP8u7G6HMl*n zGNP^3z=Y+tMQ#etPBCMu4H+m~$wHeuHT9{9ur&H*`?&!30>rZRQh&8&2Q$SEB#k$oYK!^)Mpf1q8Ke8{O4y9}%} zQE`*IlFdpw@Ky%<2#urcL1WEzHwKn3%2yLs{TqtV6^#IRoIvN0D?q6a7?m`P4Krwb##CXz)8#0IVrt G=%NLTLE1_H diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index 5b99a2b2d0ab..6c930f9a1f64 100644 --- a/codex-rs/app-server/README.md +++ b/codex-rs/app-server/README.md @@ -1709,6 +1709,8 @@ Event notifications are the server-initiated event stream for thread lifecycles, Thread realtime publishes thread-scoped timeline item lifecycle notifications for paginated threads alongside its existing realtime notifications. Completed timeline items are durably interleaved with ordinary turn items by `thread/timeline/list`. Neither surface changes `ThreadItem`, `thread/read`, `thread/resume`, or `thread/fork`; clients ignore notification methods they do not recognize. +Core records transcript segments, session boundaries, and backing-agent artifact promotions through its injected thread store, even without an app-server event listener. Presentation selection uses the same rules for every Core host. App-server translates Core's history events into the notifications below; it does not append those items again. Recording remains limited to paginated threads. A completed notification follows acceptance by the thread store, not an additional flush or power-loss durability barrier. + Each realtime item has an `id`, a `realtimeSessionId`, and one of four types: `realtimeSessionStarted`, `transcriptSegment`, `bemItemPromoted`, or `realtimeSessionClosed`. A `bemItemPromoted` item references an existing backing-agent item by `turnId` and `itemId`; its `presentation` is `wholeItem`, `inlineMarkdown`, or `inlineVisualization` with an `index`. Recoverable configuration and initialization warnings use the existing `configWarning` notification: `{ summary, details?, path?, range? }`. App-server may emit it during initialization for config parsing and related setup diagnostics, or to the requesting connection during `thread/start` when that thread's exec-policy rules fail to parse. diff --git a/codex-rs/app-server/src/bespoke_event_handling.rs b/codex-rs/app-server/src/bespoke_event_handling.rs index 7fcfeb8d3540..b16d931e5707 100644 --- a/codex-rs/app-server/src/bespoke_event_handling.rs +++ b/codex-rs/app-server/src/bespoke_event_handling.rs @@ -61,6 +61,9 @@ use codex_app_server_protocol::ThreadItem; use codex_app_server_protocol::ThreadRealtimeClosedNotification; use codex_app_server_protocol::ThreadRealtimeErrorNotification; use codex_app_server_protocol::ThreadRealtimeItemAddedNotification; +use codex_app_server_protocol::ThreadRealtimeItemCompletedNotification; +use codex_app_server_protocol::ThreadRealtimeItemStartedNotification; +use codex_app_server_protocol::ThreadRealtimeItemTranscriptDeltaNotification; use codex_app_server_protocol::ThreadRealtimeOutputAudioDeltaNotification; use codex_app_server_protocol::ThreadRealtimeSdpNotification; use codex_app_server_protocol::ThreadRealtimeStartedNotification; @@ -463,6 +466,39 @@ pub(crate) async fn apply_bespoke_event_handling( .await; } EventMsg::RealtimeConversationRealtime(event) => match event.payload { + RealtimeEvent::HistoryItemStarted(item) => { + outgoing + .send_server_notification(ServerNotification::ThreadRealtimeItemStarted( + ThreadRealtimeItemStartedNotification { + thread_id: conversation_id.to_string(), + item: item.into(), + }, + )) + .await; + } + RealtimeEvent::HistoryTranscriptDelta { item_id, delta } => { + outgoing + .send_server_notification( + ServerNotification::ThreadRealtimeItemTranscriptDelta( + ThreadRealtimeItemTranscriptDeltaNotification { + thread_id: conversation_id.to_string(), + item_id, + delta, + }, + ), + ) + .await; + } + RealtimeEvent::HistoryItemCompleted(item) => { + outgoing + .send_server_notification(ServerNotification::ThreadRealtimeItemCompleted( + ThreadRealtimeItemCompletedNotification { + thread_id: conversation_id.to_string(), + item: item.into(), + }, + )) + .await; + } RealtimeEvent::SessionUpdated { .. } => {} RealtimeEvent::InputAudioSpeechStarted(event) => { let notification = ThreadRealtimeItemAddedNotification { diff --git a/codex-rs/app-server/src/lib.rs b/codex-rs/app-server/src/lib.rs index 04991451227a..76709f869c26 100644 --- a/codex-rs/app-server/src/lib.rs +++ b/codex-rs/app-server/src/lib.rs @@ -122,8 +122,6 @@ mod models_refresh_worker; mod notification_media; mod otel_reloader; mod outgoing_message; -mod realtime_event_handling; -mod realtime_history; mod request_processors; mod request_serialization; mod server_request_error; diff --git a/codex-rs/app-server/src/realtime_event_handling.rs b/codex-rs/app-server/src/realtime_event_handling.rs deleted file mode 100644 index cc4a29e4ec4f..000000000000 --- a/codex-rs/app-server/src/realtime_event_handling.rs +++ /dev/null @@ -1,91 +0,0 @@ -use crate::outgoing_message::ThreadScopedOutgoingMessageSender; -use crate::realtime_history::RealtimeEventEffects; -use codex_app_server_protocol::ServerNotification; -use codex_app_server_protocol::ThreadRealtimeItemCompletedNotification; -use codex_app_server_protocol::ThreadRealtimeItemStartedNotification; -use codex_app_server_protocol::ThreadRealtimeItemTranscriptDeltaNotification; -use codex_core::CodexThread; -use codex_protocol::ThreadId; -use codex_protocol::realtime::RealtimeItem; -use codex_protocol::realtime::RealtimeItemContent; -use codex_rollout::RolloutItem; -use tracing::warn; - -pub(crate) async fn apply_realtime_event_effects( - conversation: &CodexThread, - outgoing: &ThreadScopedOutgoingMessageSender, - thread_id: ThreadId, - effects: RealtimeEventEffects, -) { - let thread_id = thread_id.to_string(); - - if let Some(stream) = effects.transcript_stream { - if let Some(item) = stream.started_item { - outgoing - .send_server_notification(ServerNotification::ThreadRealtimeItemStarted( - ThreadRealtimeItemStartedNotification { - thread_id: thread_id.clone(), - item: item.into(), - }, - )) - .await; - } - outgoing - .send_server_notification(ServerNotification::ThreadRealtimeItemTranscriptDelta( - ThreadRealtimeItemTranscriptDeltaNotification { - thread_id: thread_id.clone(), - item_id: stream.item_id, - delta: stream.delta, - }, - )) - .await; - } - - if let Err(error) = - persist_realtime_items(conversation, outgoing, &thread_id, effects.items).await - { - warn!(thread_id, "failed to persist realtime history: {error}"); - } -} - -pub(crate) async fn persist_realtime_items( - conversation: &CodexThread, - outgoing: &ThreadScopedOutgoingMessageSender, - thread_id: &str, - items: Vec, -) -> Result<(), String> { - if items.is_empty() { - return Ok(()); - } - conversation - .append_rollout_items( - &items - .iter() - .cloned() - .map(RolloutItem::RealtimeItem) - .collect::>(), - ) - .await - .map_err(|error| format!("failed to persist realtime history: {error}"))?; - for item in items { - if !matches!(&item.content, RealtimeItemContent::TranscriptSegment { .. }) { - outgoing - .send_server_notification(ServerNotification::ThreadRealtimeItemStarted( - ThreadRealtimeItemStartedNotification { - thread_id: thread_id.to_string(), - item: item.clone().into(), - }, - )) - .await; - } - outgoing - .send_server_notification(ServerNotification::ThreadRealtimeItemCompleted( - ThreadRealtimeItemCompletedNotification { - thread_id: thread_id.to_string(), - item: item.into(), - }, - )) - .await; - } - Ok(()) -} diff --git a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs index 7060ec4f3254..8354a4e0d37d 100644 --- a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs +++ b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs @@ -1,12 +1,8 @@ use super::*; use crate::extensions::send_thread_warning; -use crate::realtime_event_handling::apply_realtime_event_effects; -use crate::realtime_event_handling::persist_realtime_items; -use crate::realtime_history::RealtimeEventEffects; use codex_app_server_protocol::ThreadQueueChangedNotification; use codex_extension_api::ThreadIdleCause; use codex_protocol::config_types::MultiAgentMode; -use codex_protocol::protocol::ThreadHistoryMode; pub(super) const THREAD_UNLOADING_DELAY: Duration = Duration::from_secs(30 * 60); @@ -247,8 +243,6 @@ pub(super) async fn ensure_listener_task_running( ) .await; let config_snapshot = conversation.config_snapshot().await; - let realtime_history_enabled = - matches!(config_snapshot.history_mode, ThreadHistoryMode::Paginated); let thread_settings_baseline = thread_settings_from_config_snapshot(&config_snapshot); let (mut listener_command_rx, listener_generation) = { let mut thread_state = thread_state.lock().await; @@ -331,20 +325,10 @@ pub(super) async fn ensure_listener_task_running( // Track the event before emitting any typed translations // so thread-local state such as raw event opt-in stays // synchronized with the conversation. - let (raw_events_enabled, realtime_effects) = { + let raw_events_enabled = { let mut thread_state = thread_state.lock().await; thread_state.track_current_turn_event(&event.id, &event.msg); - let realtime_effects = if realtime_history_enabled - && thread_state.realtime_history.should_observe(&event.msg) - { - let active_turn_id = thread_state.active_turn_snapshot().map(|turn| turn.id); - thread_state - .realtime_history - .observe(&event.msg, active_turn_id.as_deref()) - } else { - RealtimeEventEffects::default() - }; - (thread_state.experimental_raw_events, realtime_effects) + thread_state.experimental_raw_events }; if matches!( &event.msg, @@ -362,14 +346,6 @@ pub(super) async fn ensure_listener_task_running( conversation_id, ); - apply_realtime_event_effects( - conversation.as_ref(), - &thread_outgoing, - conversation_id, - realtime_effects, - ) - .await; - apply_bespoke_event_handling( event.clone(), conversation_id, @@ -582,32 +558,6 @@ pub(super) async fn handle_thread_listener_command( .await; let _ = completion_tx.send(()); } - ThreadListenerCommand::SealRealtimeUserInput { - input, - completion_tx, - } => { - let items = thread_state - .lock() - .await - .realtime_history - .seal_user_input(&input); - let subscribed_connection_ids = thread_state_manager - .subscribed_connection_ids(conversation_id) - .await; - let thread_outgoing = ThreadScopedOutgoingMessageSender::new( - outgoing.clone(), - subscribed_connection_ids, - conversation_id, - ); - let result = persist_realtime_items( - conversation.as_ref(), - &thread_outgoing, - &conversation_id.to_string(), - items, - ) - .await; - let _ = completion_tx.send(result); - } } } diff --git a/codex-rs/app-server/src/request_processors/turn_processor.rs b/codex-rs/app-server/src/request_processors/turn_processor.rs index 129ad3b5a829..2cd4c436f8b4 100644 --- a/codex-rs/app-server/src/request_processors/turn_processor.rs +++ b/codex-rs/app-server/src/request_processors/turn_processor.rs @@ -619,10 +619,6 @@ impl TurnRequestProcessor { }, ) .await?; - if let TurnInput::UserInput { content, .. } = &input { - self.seal_realtime_transcript_before_user_input(thread_id, content) - .await?; - } let submission = thread .start_or_steer_turn( @@ -997,12 +993,12 @@ impl TurnRequestProcessor { request_id: &ConnectionRequestId, params: TurnSteerParams, ) -> Result { - let (thread_id, thread) = - self.load_thread(¶ms.thread_id) - .await - .inspect_err(|error| { - self.track_error_response(request_id, error, /*error_type*/ None); - })?; + let (_, thread) = self + .load_thread(¶ms.thread_id) + .await + .inspect_err(|error| { + self.track_error_response(request_id, error, /*error_type*/ None); + })?; self.ensure_direct_input_allowed(request_id, thread.as_ref()) .await?; @@ -1028,9 +1024,6 @@ impl TurnRequestProcessor { .collect(); let additional_context = map_additional_context(params.additional_context); - self.seal_realtime_transcript_before_user_input(thread_id, &mapped_items) - .await?; - let submission = thread .steer_turn( TurnInputRequest::new(TurnInput::UserInput { @@ -1127,37 +1120,6 @@ impl TurnRequestProcessor { Ok(TurnSteerResponse { turn_id }) } - async fn seal_realtime_transcript_before_user_input( - &self, - thread_id: ThreadId, - input: &[CoreInputItem], - ) -> Result<(), JSONRPCErrorError> { - let thread_state = self.thread_state_manager.thread_state(thread_id).await; - if !thread_state - .lock() - .await - .realtime_history - .should_seal_user_input(input) - { - return Ok(()); - } - let listener = self - .thread_state_manager - .current_listener_command_tx(thread_id) - .ok_or_else(|| internal_error("thread listener is not running"))?; - let (completion_tx, completion_rx) = tokio::sync::oneshot::channel(); - listener - .send(ThreadListenerCommand::SealRealtimeUserInput { - input: input.to_vec(), - completion_tx, - }) - .map_err(|_| internal_error("thread listener is not running"))?; - completion_rx - .await - .map_err(|_| internal_error("thread listener stopped before sealing realtime input"))? - .map_err(internal_error) - } - async fn prepare_realtime_conversation_thread( &self, request_id: &ConnectionRequestId, diff --git a/codex-rs/app-server/src/thread_state.rs b/codex-rs/app-server/src/thread_state.rs index 402783151ed2..3795638cdd36 100644 --- a/codex-rs/app-server/src/thread_state.rs +++ b/codex-rs/app-server/src/thread_state.rs @@ -1,6 +1,5 @@ use crate::outgoing_message::ConnectionId; use crate::outgoing_message::ConnectionRequestId; -use crate::realtime_history::RealtimeHistoryState; use codex_app_server_protocol::RequestId; use codex_app_server_protocol::ThreadGoal; use codex_app_server_protocol::ThreadHistoryBuilder; @@ -18,7 +17,6 @@ use codex_protocol::items::AgentMessageContent as CoreAgentMessageContent; use codex_protocol::items::TurnItem as CoreTurnItem; use codex_protocol::models::MessagePhase; use codex_protocol::protocol::EventMsg; -use codex_protocol::user_input::UserInput; use codex_rollout::RolloutItem; use codex_rollout::state_db::StateDbHandle; use codex_utils_path_uri::LegacyAppPathString; @@ -83,10 +81,6 @@ pub(crate) enum ThreadListenerCommand { request_id: RequestId, completion_tx: oneshot::Sender<()>, }, - SealRealtimeUserInput { - input: Vec, - completion_tx: oneshot::Sender>, - }, } /// Per-conversation accumulation of the latest states e.g. error message while a turn runs. @@ -110,7 +104,6 @@ pub(crate) struct ThreadState { pub(crate) cancel_tx: Option>, pub(crate) experimental_raw_events: bool, pub(crate) listener_generation: u64, - pub(crate) realtime_history: RealtimeHistoryState, last_thread_settings: Option, listener_command_tx: Option>, current_turn_history: ThreadHistoryBuilder, diff --git a/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs b/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs index 9cf86981bf08..5cc7bb3a25ea 100644 --- a/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs +++ b/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs @@ -760,10 +760,39 @@ async fn realtime_conversation_streams_timeline_items() -> Result<()> { assert_eq!(completed.item.id, started.item.id); assert_eq!(delta.delta, "hello"); assert!(matches!( - completed.item.content, + &completed.item.content, ThreadRealtimeItemContent::TranscriptSegment { text, .. } if text == "hello" )); + let closed = read_notification::( + &mut mcp, + "thread/realtime/item/completed", + ) + .await?; + let _: ThreadRealtimeClosedNotification = + read_notification(&mut mcp, "thread/realtime/closed").await?; + let request = mcp + .send_thread_timeline_list_request(ThreadTimelineListParams { + thread_id: thread.thread.id, + cursor: None, + limit: Some(100), + }) + .await?; + let page: ThreadTimelineListResponse = + timeout(DEFAULT_TIMEOUT, mcp.read_response(request)).await??; + let persisted = page + .data + .into_iter() + .filter_map(|entry| match entry { + ThreadTimelineEntry::Realtime { item, .. } => Some(item), + _ => None, + }) + .collect::>(); + assert_eq!( + persisted, + vec![session_completed.item, completed.item, closed.item] + ); + realtime_server.shutdown().await; Ok(()) } @@ -1185,6 +1214,10 @@ async fn realtime_timeline_splits_accepted_steering_and_persists_promoted_artifa "type": "response.output_text.delta", "delta": "Spoken before steering" })], + vec![json!({ + "type": "response.output_text.delta", + "delta": " and after rejected steering" + })], ])]), ) .await?; @@ -1230,6 +1263,45 @@ async fn realtime_timeline_splits_accepted_steering_and_persists_promoted_artifa ) .await?; + for (input, expected_turn_id) in [ + (Vec::new(), turn.turn.id.clone()), + ( + vec![V2UserInput::Text { + text: "Rejected steering".to_string(), + text_elements: Vec::new(), + }], + "stale-turn".to_string(), + ), + ] { + let request = harness + .mcp + .send_turn_steer_request(TurnSteerParams { + thread_id: harness.thread_id.clone(), + input, + expected_turn_id, + additional_context: None, + client_user_message_id: None, + responsesapi_client_metadata: None, + }) + .await?; + let rejected = timeout( + DEFAULT_TIMEOUT, + harness + .mcp + .read_stream_until_error_message(RequestId::Integer(request)), + ) + .await??; + assert_eq!(rejected.error.code, -32600); + } + harness + .append_text(harness.thread_id.clone(), "Continue speech after rejection") + .await?; + harness + .read_notification::( + "thread/realtime/transcript/delta", + ) + .await?; + let steering_request = harness .mcp .send_turn_steer_request(TurnSteerParams { @@ -1274,7 +1346,7 @@ async fn realtime_timeline_splits_accepted_steering_and_persists_promoted_artifa .. }, .. - } if text == "Spoken before steering" + } if text == "Spoken before steering and after rejected steering" ) }) .context("accepted steering should seal the active transcript")?; @@ -1296,16 +1368,25 @@ async fn realtime_timeline_splits_accepted_steering_and_persists_promoted_artifa }) .context("accepted steering should be included in the timeline")?; assert!(transcript_index < steering_index); - assert!(page.data.iter().any(|entry| matches!( - entry, - ThreadTimelineEntry::Realtime { - item: ThreadRealtimeItem { - content: ThreadRealtimeItemContent::BemItemPromoted { item_id, .. }, - .. - }, - .. - } if item_id == "promoted-message" - ))); + // Inline artifacts are promoted while streaming, before their final item. + assert_eq!( + page.data + .iter() + .filter_map(|entry| match entry { + ThreadTimelineEntry::Item { item, .. } + if matches!(item.as_ref(), ThreadItem::AgentMessage { id, .. } if id == "promoted-message") => Some("artifact"), + ThreadTimelineEntry::Realtime { + item: ThreadRealtimeItem { + content: ThreadRealtimeItemContent::BemItemPromoted { item_id, .. }, + .. + }, + .. + } if item_id == "promoted-message" => Some("promotion"), + _ => None, + }) + .collect::>(), + vec!["promotion", "artifact"] + ); for entry in &page.data { if let ThreadTimelineEntry::Realtime { item, .. } = entry { assert_eq!(Uuid::parse_str(&item.id)?.get_version_num(), 7); diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs index 7dd84e536613..293305ea538a 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs @@ -652,7 +652,10 @@ impl RealtimeWebsocketEvents { | RealtimeEvent::ConversationItemDone { .. } | RealtimeEvent::NoopRequested(_) | RealtimeEvent::ConversationItemAdded(_) - | RealtimeEvent::Error(_) => {} + | RealtimeEvent::Error(_) + | RealtimeEvent::HistoryItemStarted(_) + | RealtimeEvent::HistoryTranscriptDelta { .. } + | RealtimeEvent::HistoryItemCompleted(_) => {} } truncate_active_transcript(&mut active_transcript.entries); } diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 8463ae9e16e9..b4e979289267 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -11,6 +11,7 @@ mod client; mod client_common; mod realtime_context; mod realtime_conversation; +mod realtime_history; mod realtime_prompt; mod responses_metadata; mod responses_retry; diff --git a/codex-rs/core/src/realtime_conversation.rs b/codex-rs/core/src/realtime_conversation.rs index bec3a059853f..a810a48c1768 100644 --- a/codex-rs/core/src/realtime_conversation.rs +++ b/codex-rs/core/src/realtime_conversation.rs @@ -2449,7 +2449,10 @@ async fn handle_realtime_server_event( | RealtimeEvent::OutputTranscriptDelta(_) | RealtimeEvent::OutputTranscriptDone(_) | RealtimeEvent::ConversationItemAdded(_) - | RealtimeEvent::ConversationItemDone { .. } => false, + | RealtimeEvent::ConversationItemDone { .. } + | RealtimeEvent::HistoryItemStarted(_) + | RealtimeEvent::HistoryTranscriptDelta { .. } + | RealtimeEvent::HistoryItemCompleted(_) => false, }; if events_tx.send(event).await.is_err() { diff --git a/codex-rs/app-server/src/realtime_history.rs b/codex-rs/core/src/realtime_history.rs similarity index 73% rename from codex-rs/app-server/src/realtime_history.rs rename to codex-rs/core/src/realtime_history.rs index 507632f8e06d..8ba0a71f5557 100644 --- a/codex-rs/app-server/src/realtime_history.rs +++ b/codex-rs/core/src/realtime_history.rs @@ -1,12 +1,14 @@ -use codex_protocol::items::AgentMessageContent; -use codex_protocol::items::DynamicToolCallStatus; -use codex_protocol::items::McpToolCallStatus; +//! Records canonical Voice history for every Core host. Presentation rules share +//! transcript boundaries and deduplication state with the history reducer. +//! The reducer selects when its effects are persisted relative to their source event. + +mod presentation; + use codex_protocol::items::TurnItem; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::RealtimeEvent; use codex_protocol::protocol::RealtimeTranscriptDelta; use codex_protocol::protocol::RealtimeTranscriptDone; -use codex_protocol::protocol::SubAgentActivityKind; use codex_protocol::realtime::BemItemPresentation; use codex_protocol::realtime::RealtimeItem; use codex_protocol::realtime::RealtimeItemContent; @@ -18,12 +20,6 @@ use std::collections::HashSet; use std::collections::VecDeque; use uuid::Uuid; -const INLINE_MARKDOWN_DIRECTIVE: &str = "::codex-realtime-inline{}"; -const INLINE_VISUALIZATION_DIRECTIVE: &str = "::codex-inline-vis{"; -const VISUALIZE_DIRECTIVE: &str = "visualize{"; -const BACKTICK_FENCE: &str = "```"; -const TILDE_FENCE: &str = "~~~"; - #[derive(Debug, Clone, PartialEq, Eq)] struct ActiveSegment { session_id: String, @@ -75,10 +71,18 @@ pub(crate) struct RealtimeTranscriptStream { #[derive(Debug, Default)] pub(crate) struct RealtimeEventEffects { + pub(crate) order: RealtimeEventOrder, pub(crate) items: Vec, pub(crate) transcript_stream: Option, } +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub(crate) enum RealtimeEventOrder { + BeforeEvent, + #[default] + AfterEvent, +} + #[derive(Clone, Copy)] enum Continuation { Continue, @@ -86,11 +90,12 @@ enum Continuation { } /// Retains only live session state; durable history is served by the rollout index. -#[derive(Debug, Default)] +#[derive(Default)] pub(crate) struct RealtimeHistoryState { active_session_id: Option, active_segments: ActiveTranscriptSegments, streaming_agent_message: Option, + active_turn_id: Option, realtime_session_by_bem_turn: HashMap, promoted_bem_presentation_keys: HashSet, pending_handoffs: VecDeque, @@ -98,7 +103,7 @@ pub(crate) struct RealtimeHistoryState { } impl RealtimeHistoryState { - pub(crate) fn should_seal_user_input(&self, input: &[UserInput]) -> bool { + fn should_seal_user_input(&self, input: &[UserInput]) -> bool { self.active_session_id.is_some() && [&self.active_segments.user, &self.active_segments.assistant] .into_iter() @@ -111,7 +116,7 @@ impl RealtimeHistoryState { }) } - pub(crate) fn seal_user_input(&mut self, input: &[UserInput]) -> Vec { + fn seal_user_input(&mut self, input: &[UserInput]) -> Vec { if !self.should_seal_user_input(input) { return Vec::new(); } @@ -121,17 +126,22 @@ impl RealtimeHistoryState { } pub(crate) fn should_observe(&self, event: &EventMsg) -> bool { - matches!(event, EventMsg::RealtimeConversationStarted(_)) - || (self.active_session_id.is_some() - && matches!( - event, - EventMsg::RealtimeConversationRealtime(_) - | EventMsg::RealtimeConversationClosed(_) - | EventMsg::TurnStarted(_) - | EventMsg::ItemStarted(_) - | EventMsg::ItemCompleted(_) - | EventMsg::AgentMessageContentDelta(_) - )) + matches!( + event, + EventMsg::RealtimeConversationStarted(_) + | EventMsg::TurnStarted(_) + | EventMsg::TurnComplete(_) + | EventMsg::TurnAborted(_) + ) || (self.active_session_id.is_some() + && matches!( + event, + EventMsg::RealtimeConversationRealtime(_) + | EventMsg::RealtimeConversationClosed(_) + | EventMsg::TurnStarted(_) + | EventMsg::ItemStarted(_) + | EventMsg::ItemCompleted(_) + | EventMsg::AgentMessageContentDelta(_) + )) || match event { EventMsg::ItemStarted(event) => self .realtime_session_by_bem_turn @@ -142,16 +152,29 @@ impl RealtimeHistoryState { EventMsg::AgentMessageContentDelta(event) => self .realtime_session_by_bem_turn .contains_key(&event.turn_id), - EventMsg::TurnStarted(_) => !self.pending_handoffs.is_empty(), _ => false, } } - pub(crate) fn observe( - &mut self, - event: &EventMsg, - active_turn_id: Option<&str>, - ) -> RealtimeEventEffects { + pub(crate) fn observe(&mut self, event: &EventMsg) -> RealtimeEventEffects { + match event { + EventMsg::TurnStarted(event) => self.active_turn_id = Some(event.turn_id.clone()), + EventMsg::TurnComplete(event) + if self.active_turn_id.as_deref() == Some(event.turn_id.as_str()) => + { + self.active_turn_id = None; + } + EventMsg::TurnAborted(event) + if event.turn_id.is_none() || event.turn_id == self.active_turn_id => + { + self.active_turn_id = None; + } + _ => {} + } + if !self.should_observe(event) { + return RealtimeEventEffects::default(); + } + let mut order = RealtimeEventOrder::AfterEvent; let mut items = Vec::new(); let mut transcript_stream = None; match event { @@ -170,7 +193,7 @@ impl RealtimeHistoryState { content: RealtimeItemContent::RealtimeSessionStarted, }); } - if let Some(turn_id) = active_turn_id { + if let Some(turn_id) = &self.active_turn_id { self.realtime_session_by_bem_turn .insert(turn_id.to_string(), session_id); } @@ -211,6 +234,7 @@ impl RealtimeHistoryState { } EventMsg::ItemStarted(event) => { if let TurnItem::UserMessage(item) = &event.item { + order = RealtimeEventOrder::BeforeEvent; items.extend(self.seal_user_input(&item.content)); } self.observe_item( @@ -222,6 +246,7 @@ impl RealtimeHistoryState { } EventMsg::ItemCompleted(event) => { if let TurnItem::UserMessage(item) = &event.item { + order = RealtimeEventOrder::BeforeEvent; items.extend(self.seal_user_input(&item.content)); } self.observe_item( @@ -273,121 +298,12 @@ impl RealtimeHistoryState { _ => {} } RealtimeEventEffects { + order, items, transcript_stream, } } - fn observe_item( - &mut self, - items: &mut Vec, - turn_id: &str, - item: &TurnItem, - completed: bool, - ) { - match item { - TurnItem::AgentMessage(message) => { - let text = message - .content - .iter() - .map(|content| match content { - AgentMessageContent::Text { text } => text.as_str(), - }) - .collect::(); - if !completed { - self.streaming_agent_message = Some(StreamingAgentMessage { - item_id: message.id.clone(), - text: text.clone(), - }); - } - self.observe_assistant_message(items, turn_id, &message.id, &text); - } - TurnItem::ImageGeneration(image) => { - self.add_promotion(items, turn_id, &image.id, BemItemPresentation::WholeItem); - } - TurnItem::Extension(extension) - if serde_json::to_value(extension) - .is_ok_and(|item| item["kind"] == "image_gen.generation") => - { - self.add_promotion( - items, - turn_id, - extension.id(), - BemItemPresentation::WholeItem, - ); - } - TurnItem::SubAgentActivity(activity) - if completed && activity.kind == SubAgentActivityKind::Started => - { - self.add_promotion(items, turn_id, &activity.id, BemItemPresentation::WholeItem); - } - TurnItem::DynamicToolCall(call) - if self.active_session_id.is_some() - && completed - && call.status == DynamicToolCallStatus::Completed - && call.success == Some(true) => - { - self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem); - } - TurnItem::McpToolCall(call) - if self.active_session_id.is_some() - && completed - && call.server == "codex_app" - && call.status == McpToolCallStatus::Completed => - { - self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem); - } - _ => {} - } - } - - fn observe_assistant_message( - &mut self, - items: &mut Vec, - turn_id: &str, - item_id: &str, - text: &str, - ) { - let mut lines = text.trim_start().lines(); - let mut first = lines.next().unwrap_or_default(); - if first.starts_with('[') - && let Some((_, content)) = first.split_once(']') - { - first = content.trim_start(); - if first.is_empty() { - first = lines.next().unwrap_or_default(); - } - } - if first == INLINE_MARKDOWN_DIRECTIVE && text.contains('\n') { - self.add_promotion(items, turn_id, item_id, BemItemPresentation::InlineMarkdown); - return; - } - - let mut in_fence = false; - let mut visualization_index = 0; - for line in text.lines() { - let trimmed = line.trim_start(); - if trimmed.starts_with(BACKTICK_FENCE) || trimmed.starts_with(TILDE_FENCE) { - in_fence = !in_fence; - continue; - } - if !in_fence - && (trimmed.starts_with(INLINE_VISUALIZATION_DIRECTIVE) - || trimmed.starts_with(VISUALIZE_DIRECTIVE)) - { - self.add_promotion( - items, - turn_id, - item_id, - BemItemPresentation::InlineVisualization { - index: visualization_index, - }, - ); - visualization_index += 1; - } - } - } - fn add_promotion( &mut self, items: &mut Vec, diff --git a/codex-rs/core/src/realtime_history/presentation.rs b/codex-rs/core/src/realtime_history/presentation.rs new file mode 100644 index 000000000000..f21c0d02b523 --- /dev/null +++ b/codex-rs/core/src/realtime_history/presentation.rs @@ -0,0 +1,129 @@ +//! Selects backing-agent items for the canonical Voice timeline using shared rules. + +use super::RealtimeHistoryState; +use super::StreamingAgentMessage; +use codex_protocol::items::AgentMessageContent; +use codex_protocol::items::DynamicToolCallStatus; +use codex_protocol::items::McpToolCallStatus; +use codex_protocol::items::TurnItem; +use codex_protocol::protocol::SubAgentActivityKind; +use codex_protocol::realtime::BemItemPresentation; +use codex_protocol::realtime::RealtimeItem; + +const INLINE_MARKDOWN_DIRECTIVE: &str = "::codex-realtime-inline{}"; +const INLINE_VISUALIZATION_DIRECTIVE: &str = "::codex-inline-vis{"; +const VISUALIZE_DIRECTIVE: &str = "visualize{"; +const BACKTICK_FENCE: &str = "```"; +const TILDE_FENCE: &str = "~~~"; + +impl RealtimeHistoryState { + pub(super) fn observe_item( + &mut self, + items: &mut Vec, + turn_id: &str, + item: &TurnItem, + completed: bool, + ) { + match item { + TurnItem::AgentMessage(message) => { + let text = message + .content + .iter() + .map(|content| match content { + AgentMessageContent::Text { text } => text.as_str(), + }) + .collect::(); + if !completed { + self.streaming_agent_message = Some(StreamingAgentMessage { + item_id: message.id.clone(), + text: text.clone(), + }); + } + self.observe_assistant_message(items, turn_id, &message.id, &text); + } + TurnItem::ImageGeneration(image) => { + self.add_promotion(items, turn_id, &image.id, BemItemPresentation::WholeItem); + } + TurnItem::Extension(extension) + if serde_json::to_value(extension) + .is_ok_and(|item| item["kind"] == "image_gen.generation") => + { + self.add_promotion( + items, + turn_id, + extension.id(), + BemItemPresentation::WholeItem, + ); + } + TurnItem::SubAgentActivity(activity) + if completed && activity.kind == SubAgentActivityKind::Started => + { + self.add_promotion(items, turn_id, &activity.id, BemItemPresentation::WholeItem); + } + TurnItem::DynamicToolCall(call) + if self.active_session_id.is_some() + && completed + && call.status == DynamicToolCallStatus::Completed + && call.success == Some(true) => + { + self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem); + } + TurnItem::McpToolCall(call) + if self.active_session_id.is_some() + && completed + && call.server == "codex_app" + && call.status == McpToolCallStatus::Completed => + { + self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem); + } + _ => {} + } + } + + pub(super) fn observe_assistant_message( + &mut self, + items: &mut Vec, + turn_id: &str, + item_id: &str, + text: &str, + ) { + let mut lines = text.trim_start().lines(); + let mut first = lines.next().unwrap_or_default(); + if first.starts_with('[') + && let Some((_, content)) = first.split_once(']') + { + first = content.trim_start(); + if first.is_empty() { + first = lines.next().unwrap_or_default(); + } + } + if first == INLINE_MARKDOWN_DIRECTIVE && text.contains('\n') { + self.add_promotion(items, turn_id, item_id, BemItemPresentation::InlineMarkdown); + return; + } + + let mut in_fence = false; + let mut visualization_index = 0; + for line in text.lines() { + let trimmed = line.trim_start(); + if trimmed.starts_with(BACKTICK_FENCE) || trimmed.starts_with(TILDE_FENCE) { + in_fence = !in_fence; + continue; + } + if !in_fence + && (trimmed.starts_with(INLINE_VISUALIZATION_DIRECTIVE) + || trimmed.starts_with(VISUALIZE_DIRECTIVE)) + { + self.add_promotion( + items, + turn_id, + item_id, + BemItemPresentation::InlineVisualization { + index: visualization_index, + }, + ); + visualization_index += 1; + } + } + } +} diff --git a/codex-rs/app-server/src/realtime_history_tests.rs b/codex-rs/core/src/realtime_history_tests.rs similarity index 65% rename from codex-rs/app-server/src/realtime_history_tests.rs rename to codex-rs/core/src/realtime_history_tests.rs index ad16bcff4590..7e0afae7b554 100644 --- a/codex-rs/app-server/src/realtime_history_tests.rs +++ b/codex-rs/core/src/realtime_history_tests.rs @@ -1,11 +1,15 @@ use super::*; use codex_protocol::AgentPath; use codex_protocol::ThreadId; +use codex_protocol::items::AgentMessageContent; use codex_protocol::items::AgentMessageItem; use codex_protocol::items::DynamicToolCallItem; +use codex_protocol::items::DynamicToolCallStatus; use codex_protocol::items::ImageGenerationItem; use codex_protocol::items::McpToolCallItem; +use codex_protocol::items::McpToolCallStatus; use codex_protocol::items::SubAgentActivityItem; +use codex_protocol::items::UserMessageItem; use codex_protocol::protocol::AgentMessageContentDeltaEvent; use codex_protocol::protocol::ItemCompletedEvent; use codex_protocol::protocol::ItemStartedEvent; @@ -13,18 +17,30 @@ use codex_protocol::protocol::RealtimeConversationClosedEvent; use codex_protocol::protocol::RealtimeConversationRealtimeEvent; use codex_protocol::protocol::RealtimeConversationStartedEvent; use codex_protocol::protocol::RealtimeConversationVersion; +use codex_protocol::protocol::SubAgentActivityKind; +use codex_protocol::protocol::TurnAbortReason; +use codex_protocol::protocol::TurnAbortedEvent; use pretty_assertions::assert_eq; +use test_case::test_case; use uuid::Uuid; fn started_state() -> RealtimeHistoryState { let mut state = RealtimeHistoryState::default(); - state.observe( - &EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { + state.observe(&EventMsg::TurnStarted( + codex_protocol::protocol::TurnStartedEvent { + turn_id: "turn-1".to_string(), + trace_id: None, + started_at: None, + model_context_window: None, + collaboration_mode_kind: Default::default(), + }, + )); + state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { realtime_session_id: Some("voice-1".to_string()), version: RealtimeConversationVersion::V2, - }), - Some("turn-1"), - ); + }, + )); state } @@ -32,10 +48,9 @@ fn observe_realtime( state: &mut RealtimeHistoryState, payload: RealtimeEvent, ) -> RealtimeEventEffects { - state.observe( - &EventMsg::RealtimeConversationRealtime(RealtimeConversationRealtimeEvent { payload }), - /*active_turn_id*/ None, - ) + state.observe(&EventMsg::RealtimeConversationRealtime( + RealtimeConversationRealtimeEvent { payload }, + )) } fn assistant_delta(item_id: &str, delta: &str) -> EventMsg { @@ -102,20 +117,96 @@ fn contents(effects: RealtimeEventEffects) -> Vec { .collect() } +#[test_case(Some("turn-1"), None; "matching_turn")] +#[test_case(None, None; "without_turn_id")] +#[test_case(Some("another-turn"), Some("voice-2"); "unrelated_turn")] +fn interrupted_turn_is_not_associated_with_a_new_voice_session( + aborted_turn_id: Option<&str>, + expected_session_id: Option<&str>, +) { + let mut state = RealtimeHistoryState::default(); + state.observe(&EventMsg::TurnStarted( + codex_protocol::protocol::TurnStartedEvent { + turn_id: "turn-1".to_string(), + trace_id: None, + started_at: None, + model_context_window: None, + collaboration_mode_kind: Default::default(), + }, + )); + let aborted = EventMsg::TurnAborted(TurnAbortedEvent { + turn_id: aborted_turn_id.map(str::to_string), + reason: TurnAbortReason::Interrupted, + started_at: None, + completed_at: None, + duration_ms: None, + }); + assert!(state.should_observe(&aborted)); + state.observe(&aborted); + state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { + realtime_session_id: Some("voice-2".to_string()), + version: RealtimeConversationVersion::V2, + }, + )); + + let promoted = state + .observe(&completed_item(dynamic_tool_call( + "late-tool", + DynamicToolCallStatus::Completed, + ))) + .items; + assert_eq!( + promoted + .iter() + .map(|item| item.realtime_session_id.as_str()) + .collect::>(), + expected_session_id.into_iter().collect::>() + ); +} + +#[test] +fn interrupted_turn_keeps_its_existing_voice_session_for_late_artifacts() { + let mut state = started_state(); + state.observe(&EventMsg::TurnAborted(TurnAbortedEvent { + turn_id: Some("turn-1".to_string()), + reason: TurnAbortReason::Interrupted, + started_at: None, + completed_at: None, + duration_ms: None, + })); + state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { + realtime_session_id: Some("voice-2".to_string()), + version: RealtimeConversationVersion::V2, + }, + )); + let promoted = state + .observe(&completed_item(dynamic_tool_call( + "late-tool", + DynamicToolCallStatus::Completed, + ))) + .items; + assert_eq!( + promoted + .iter() + .map(|item| item.realtime_session_id.as_str()) + .collect::>(), + vec!["voice-1"] + ); +} + #[test] fn promotes_backing_agent_artifacts_once_without_a_client_request() { let mut state = started_state(); let first_delta = assistant_delta("message-1", "[analysis] ::codex-realtime-inline{}"); - assert!( - state - .observe(&first_delta, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&first_delta).items.is_empty()); let second_delta = assistant_delta("message-1", "\nVisible explanation"); + let effects = state.observe(&second_delta); + assert_eq!(effects.order, RealtimeEventOrder::AfterEvent); assert_eq!( - contents(state.observe(&second_delta, /*active_turn_id*/ None)), + contents(effects), vec![RealtimeItemContent::BemItemPromoted { turn_id: "turn-1".to_string(), item_id: "message-1".to_string(), @@ -132,16 +223,11 @@ fn promotes_backing_agent_artifacts_once_without_a_client_request() { memory_citation: None, delivery: None, })); - assert!( - state - .observe(&completed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&completed).items.is_empty()); let next_message = assistant_delta("message-2", "::codex-realtime-inline{}\nNext message"); assert_eq!( - contents(state.observe(&next_message, /*active_turn_id*/ None)), + contents(state.observe(&next_message)), vec![RealtimeItemContent::BemItemPromoted { turn_id: "turn-1".to_string(), item_id: "message-2".to_string(), @@ -169,7 +255,7 @@ fn promotes_backing_agent_artifacts_once_without_a_client_request() { })); for (event, item_id) in [(image, "image-1"), (subagent, "subagent-1")] { assert_eq!( - contents(state.observe(&event, /*active_turn_id*/ None)), + contents(state.observe(&event)), vec![RealtimeItemContent::BemItemPromoted { turn_id: "turn-1".to_string(), item_id: item_id.to_string(), @@ -177,6 +263,15 @@ fn promotes_backing_agent_artifacts_once_without_a_client_request() { }] ); } + state.observe(&EventMsg::RealtimeConversationClosed( + RealtimeConversationClosedEvent { reason: None }, + )); + let late = state.observe(&assistant_delta( + "late-artifact", + "::codex-realtime-inline{}\npresent later", + )); + assert_eq!(late.items.len(), 1); + assert_eq!(late.items[0].realtime_session_id, "voice-1"); } #[test] @@ -185,7 +280,7 @@ fn promotes_distinct_visualizations_once_and_ignores_markdown_fences() { let text = "```\n::codex-inline-vis{file=hidden}\n```\n~~~\nvisualize{file=also-hidden}\n~~~\n::codex-inline-vis{file=first}\nvisualize{file=second}"; let message = assistant_delta("message-1", text); assert_eq!( - contents(state.observe(&message, /*active_turn_id*/ None)), + contents(state.observe(&message)), (0..2) .map(|index| RealtimeItemContent::BemItemPromoted { turn_id: "turn-1".to_string(), @@ -204,12 +299,7 @@ fn promotes_distinct_visualizations_once_and_ignores_markdown_fences() { memory_citation: None, delivery: None, })); - assert!( - state - .observe(&completed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&completed).items.is_empty()); } #[test] @@ -302,12 +392,9 @@ fn preserves_first_transcript_activity_order_when_both_speakers_are_active() { ); } let items = state - .observe( - &EventMsg::RealtimeConversationClosed(RealtimeConversationClosedEvent { - reason: None, - }), - /*active_turn_id*/ None, - ) + .observe(&EventMsg::RealtimeConversationClosed( + RealtimeConversationClosedEvent { reason: None }, + )) .items; let observed_roles = items .iter() @@ -333,13 +420,7 @@ fn does_not_replay_a_final_transcript_after_a_promotion_split() { "successful-tool", DynamicToolCallStatus::Completed, )); - assert_eq!( - state - .observe(&promotion, /*active_turn_id*/ None) - .items - .len(), - 2 - ); + assert_eq!(state.observe(&promotion).items.len(), 2); let done = observe_realtime( &mut state, @@ -356,17 +437,15 @@ fn generates_distinct_boundary_ids_for_reused_realtime_sessions() { let mut state = RealtimeHistoryState::default(); let mut boundary_ids = Vec::new(); for _ in 0..2 { - let started = state.observe( - &EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { + let started = state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { realtime_session_id: Some("reused-session".to_string()), version: RealtimeConversationVersion::V2, - }), - /*active_turn_id*/ None, - ); - let closed = state.observe( - &EventMsg::RealtimeConversationClosed(RealtimeConversationClosedEvent { reason: None }), - /*active_turn_id*/ None, - ); + }, + )); + let closed = state.observe(&EventMsg::RealtimeConversationClosed( + RealtimeConversationClosedEvent { reason: None }, + )); for item in started.items.into_iter().chain(closed.items) { assert_eq!(item.realtime_session_id, "reused-session"); assert_eq!(Uuid::parse_str(&item.id).unwrap().get_version_num(), 7); @@ -378,13 +457,12 @@ fn generates_distinct_boundary_ids_for_reused_realtime_sessions() { .collect::>(); assert_eq!(unique_ids.len(), 4); - let fallback = state.observe( - &EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { + let fallback = state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { realtime_session_id: None, version: RealtimeConversationVersion::V2, - }), - /*active_turn_id*/ None, - ); + }, + )); assert_eq!(fallback.items.len(), 1); assert_eq!( Uuid::parse_str(&fallback.items[0].realtime_session_id) @@ -410,10 +488,11 @@ fn promotes_successful_codex_app_mcp_calls_and_splits_active_transcripts() { ] { assert!( state - .observe( - &completed_item(mcp_tool_call("not-promoted", server, status)), - /*active_turn_id*/ None, - ) + .observe(&completed_item(mcp_tool_call( + "not-promoted", + server, + status + )),) .items .is_empty() ); @@ -424,7 +503,7 @@ fn promotes_successful_codex_app_mcp_calls_and_splits_active_transcripts() { McpToolCallStatus::Completed, )); assert_eq!( - contents(state.observe(&completed, /*active_turn_id*/ None)), + contents(state.observe(&completed)), vec![ RealtimeItemContent::TranscriptSegment { role: RealtimeTranscriptRole::Assistant, @@ -437,12 +516,7 @@ fn promotes_successful_codex_app_mcp_calls_and_splits_active_transcripts() { }, ] ); - assert!( - state - .observe(&completed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&completed).items.is_empty()); } #[test] @@ -461,18 +535,15 @@ fn promotes_successful_dynamic_tools_and_splits_active_transcripts() { "failed-tool", DynamicToolCallStatus::Failed, )); - assert!( - state - .observe(&failed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&failed).items.is_empty()); let completed = completed_item(dynamic_tool_call( "successful-tool", DynamicToolCallStatus::Completed, )); - let items = state.observe(&completed, /*active_turn_id*/ None).items; + let effects = state.observe(&completed); + assert_eq!(effects.order, RealtimeEventOrder::AfterEvent); + let items = effects.items; assert_eq!(items.len(), 2); assert_eq!( items[0], @@ -505,37 +576,89 @@ fn promotes_successful_dynamic_tools_and_splits_active_transcripts() { .expect("second transcript stream"); assert_ne!(first.item_id, second.item_id); assert!(second.started_item.is_some()); - assert!( - state - .observe(&completed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&completed).items.is_empty()); - let closed = state.observe( - &EventMsg::RealtimeConversationClosed(RealtimeConversationClosedEvent { reason: None }), - /*active_turn_id*/ None, - ); + let closed = state.observe(&EventMsg::RealtimeConversationClosed( + RealtimeConversationClosedEvent { reason: None }, + )); assert_eq!(closed.items[0].id, second.item_id); let late = completed_item(dynamic_tool_call( "late-tool", DynamicToolCallStatus::Completed, )); + assert!(state.observe(&late).items.is_empty()); + + state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { + realtime_session_id: Some("voice-2".to_string()), + version: RealtimeConversationVersion::V2, + }, + )); + let resumed = state.observe(&late).items; + assert_eq!(resumed.len(), 1); + assert_eq!(resumed[0].realtime_session_id, "voice-2"); +} + +#[test_case(|item| EventMsg::ItemStarted(ItemStartedEvent { + thread_id: ThreadId::new(), + turn_id: "turn-1".to_string(), + item, + started_at_ms: 0, +}); "started")] +#[test_case(completed_item; "completed")] +fn typed_input_seals_both_roles_but_realtime_delegation_does_not( + user_event: fn(TurnItem) -> EventMsg, +) { + let mut state = started_state(); + for event in [ + RealtimeEvent::OutputTranscriptDelta(RealtimeTranscriptDelta { + delta: "assistant".to_string(), + }), + RealtimeEvent::InputTranscriptDelta(RealtimeTranscriptDelta { + delta: "user".to_string(), + }), + ] { + observe_realtime(&mut state, event); + } assert!( state - .observe(&late, /*active_turn_id*/ None) + .observe(&user_event(TurnItem::UserMessage(UserMessageItem::new(&[ + UserInput::Text { + text: "delegated".to_string(), + text_elements: Vec::new(), + } + ])))) .items .is_empty() ); - - state.observe( - &EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { - realtime_session_id: Some("voice-2".to_string()), - version: RealtimeConversationVersion::V2, - }), - Some("turn-1"), + let effects = state.observe(&user_event(TurnItem::UserMessage(UserMessageItem::new(&[ + UserInput::Text { + text: "typed steering".to_string(), + text_elements: Vec::new(), + }, + ])))); + assert_eq!(effects.order, RealtimeEventOrder::BeforeEvent); + assert_eq!( + contents(effects), + vec![ + RealtimeItemContent::TranscriptSegment { + role: RealtimeTranscriptRole::Assistant, + text: "assistant".to_string() + }, + RealtimeItemContent::TranscriptSegment { + role: RealtimeTranscriptRole::User, + text: "user".to_string() + }, + ] ); - let resumed = state.observe(&late, /*active_turn_id*/ None).items; - assert_eq!(resumed.len(), 1); - assert_eq!(resumed[0].realtime_session_id, "voice-2"); + for event in [ + RealtimeEvent::OutputTranscriptDone(RealtimeTranscriptDone { + text: "assistant revised".to_string(), + }), + RealtimeEvent::InputTranscriptDone(RealtimeTranscriptDone { + text: "user revised".to_string(), + }), + ] { + assert!(observe_realtime(&mut state, event).items.is_empty()); + } } diff --git a/codex-rs/core/src/session/mod.rs b/codex-rs/core/src/session/mod.rs index f7ab1271a0cb..aa0bfe12cc31 100644 --- a/codex-rs/core/src/session/mod.rs +++ b/codex-rs/core/src/session/mod.rs @@ -40,6 +40,7 @@ use crate::image_preparation::prepare_response_items as prepare_image_response_i use crate::image_preparation::unified_image_budget_enabled; use crate::parse_turn_item; use crate::realtime_conversation::RealtimeConversationManager; +use crate::realtime_history::RealtimeEventOrder; use crate::session::step_context::StepContext; use crate::session::step_settings::ResolvedStepSettings; use crate::session::step_settings::StepSettings; @@ -227,6 +228,7 @@ mod mcp_prewarm; mod mcp_refresh; mod mcp_runtime; pub(crate) mod multi_agents; +mod realtime_history; mod review; mod rollout_budget; mod rollout_reconstruction; @@ -2339,7 +2341,32 @@ impl Session { } async fn send_event_raw_with_persistence(&self, event: Event, persist: bool) { + // Keep realtime reduction, canonical append, and delivery in the same order. + // This lock must not acquire SessionState or ActiveTurn: event producers can + // already hold those locks. Host presentation policies are synchronous. + let mut realtime_history = match &self.realtime_history { + Some(history) => { + let history = history.lock().await; + history.should_observe(&event.msg).then_some(history) + } + None => None, + }; self.services.mcp_runtime.observe_event(&event.msg); + let (before_event, after_event) = match realtime_history.as_mut() { + Some(history) => { + let effects = history.observe(&event.msg); + match effects.order { + RealtimeEventOrder::BeforeEvent => (Some(effects), None), + RealtimeEventOrder::AfterEvent => (None, Some(effects)), + } + } + None => (None, None), + }; + if let Some(effects) = before_event + && let Err(error) = self.send_realtime_history_effects(&event.id, effects).await + { + warn!("failed to persist realtime history: {error}"); + } // Persist the event into rollout storage; the store applies its persistence policy. if persist { let rollout_items = vec![RolloutItem::EventMsg(event.msg.clone())]; @@ -2348,6 +2375,11 @@ impl Session { self.services .rollout_thread_trace .record_protocol_event(&event.msg); + if let Some(effects) = after_event + && let Err(error) = self.send_realtime_history_effects(&event.id, effects).await + { + warn!("failed to persist realtime history: {error}"); + } self.deliver_event_raw(event).await; } diff --git a/codex-rs/core/src/session/realtime_history.rs b/codex-rs/core/src/session/realtime_history.rs new file mode 100644 index 000000000000..5b1acf4873bb --- /dev/null +++ b/codex-rs/core/src/session/realtime_history.rs @@ -0,0 +1,56 @@ +use super::session::Session; +use crate::realtime_history::RealtimeEventEffects; +use codex_history::RolloutItem; +use codex_protocol::protocol::Event; +use codex_protocol::protocol::EventMsg; +use codex_protocol::protocol::RealtimeConversationRealtimeEvent; +use codex_protocol::protocol::RealtimeEvent; +use codex_protocol::realtime::RealtimeItemContent; + +impl Session { + pub(super) async fn send_realtime_history_effects( + &self, + submission_id: &str, + effects: RealtimeEventEffects, + ) -> anyhow::Result<()> { + let event = |payload| Event { + id: submission_id.to_string(), + msg: EventMsg::RealtimeConversationRealtime(RealtimeConversationRealtimeEvent { + payload, + }), + }; + if let Some(stream) = effects.transcript_stream { + if let Some(item) = stream.started_item { + self.deliver_event_raw(event(RealtimeEvent::HistoryItemStarted(item))) + .await; + } + self.deliver_event_raw(event(RealtimeEvent::HistoryTranscriptDelta { + item_id: stream.item_id, + delta: stream.delta, + })) + .await; + } + if effects.items.is_empty() { + return Ok(()); + } + self.live_thread_for_persistence("append realtime history")? + .append_items( + &effects + .items + .iter() + .cloned() + .map(RolloutItem::RealtimeItem) + .collect::>(), + ) + .await?; + for item in effects.items { + if !matches!(&item.content, RealtimeItemContent::TranscriptSegment { .. }) { + self.deliver_event_raw(event(RealtimeEvent::HistoryItemStarted(item.clone()))) + .await; + } + self.deliver_event_raw(event(RealtimeEvent::HistoryItemCompleted(item))) + .await; + } + Ok(()) + } +} diff --git a/codex-rs/core/src/session/session.rs b/codex-rs/core/src/session/session.rs index 3b3da7440ba8..63530cc31309 100644 --- a/codex-rs/core/src/session/session.rs +++ b/codex-rs/core/src/session/session.rs @@ -64,6 +64,7 @@ pub(crate) struct Session { pub(super) mcp_prewarm_shutdown: CancellationToken, pub(super) mcp_prewarm_task: std::sync::Mutex>>, pub(crate) conversation: Arc, + pub(crate) realtime_history: Option>, pub(crate) active_turn: Mutex>, pub(crate) async_hook_results: async_channel::Receiver, pub(crate) input_queue: InputQueue, @@ -1492,6 +1493,9 @@ impl Session { mcp_prewarm_shutdown: CancellationToken::new(), mcp_prewarm_task: std::sync::Mutex::new(None), conversation: Arc::new(RealtimeConversationManager::new()), + realtime_history: (session_configuration.history_mode == ThreadHistoryMode::Paginated + && services.live_thread.is_some()) + .then(|| Mutex::new(Default::default())), active_turn: Mutex::new(None), async_hook_results, input_queue: InputQueue::new(), diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index b8d97c87e526..268de9339131 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -6612,6 +6612,7 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { mcp_prewarm_shutdown: CancellationToken::new(), mcp_prewarm_task: std::sync::Mutex::new(None), conversation: Arc::new(RealtimeConversationManager::new()), + realtime_history: None, active_turn: Mutex::new(None), async_hook_results, input_queue: super::input_queue::InputQueue::new(), @@ -8892,6 +8893,7 @@ where mcp_prewarm_shutdown: CancellationToken::new(), mcp_prewarm_task: std::sync::Mutex::new(None), conversation: Arc::new(RealtimeConversationManager::new()), + realtime_history: None, active_turn: Mutex::new(None), async_hook_results, input_queue: super::input_queue::InputQueue::new(), diff --git a/codex-rs/core/src/session/turn_input_tests.rs b/codex-rs/core/src/session/turn_input_tests.rs index 3ef6797c065b..c9810df1f4cb 100644 --- a/codex-rs/core/src/session/turn_input_tests.rs +++ b/codex-rs/core/src/session/turn_input_tests.rs @@ -1,6 +1,7 @@ use super::*; use crate::config::Constrained; use crate::session::step_settings::StepSettingsUpdate; +use crate::session::tests::make_session_and_context; use crate::session::tests::make_session_and_context_with_rx; use crate::session::turn_context::TurnContext; use crate::state::TaskKind; @@ -107,6 +108,65 @@ async fn submit_steer_only( .expect("steer-only submission should be valid") } +#[tokio::test] +#[expect( + clippy::await_holding_invalid_type, + reason = "simulate an in-flight realtime append while checking input admission" +)] +async fn steering_does_not_wait_for_realtime_history() { + let (mut session, turn_context) = make_session_and_context().await; + session.realtime_history = Some(tokio::sync::Mutex::new(Default::default())); + let session = Arc::new(session); + let turn_context = Arc::new(turn_context); + session + .spawn_task( + Arc::clone(&turn_context), + Vec::new(), + NeverEndingTask { + kind: TaskKind::Regular, + listen_to_cancellation_token: true, + }, + ) + .await; + + let history = session + .realtime_history + .as_ref() + .expect("realtime history") + .lock() + .await; + for mode in [ + TurnInputMode::StartOrSteer, + TurnInputMode::Steer { + expected_turn_id: turn_context.sub_id.clone(), + }, + ] { + let submission = tokio::time::timeout( + std::time::Duration::from_secs(5), + handle( + &session, + TurnInputRequest::user_input(vec![UserInput::Text { + text: "steer without waiting for persistence".to_string(), + text_elements: Vec::new(), + }]), + mode, + "steer-submission".to_string(), + ), + ) + .await + .expect("steering must not wait for the realtime recorder") + .expect("steering should succeed"); + assert_eq!( + submission, + TurnInputSubmission::Steered { + turn_id: turn_context.sub_id.clone() + } + ); + } + drop(history); + session.abort_all_tasks(TurnAbortReason::Interrupted).await; +} + #[tokio::test] async fn accepted_input_applies_thread_settings() { let (session, turn_context, _rx) = make_session_and_context_with_rx().await; diff --git a/codex-rs/core/tests/suite/realtime_conversation.rs b/codex-rs/core/tests/suite/realtime_conversation.rs index de16d730ed57..b891fcfd6436 100644 --- a/codex-rs/core/tests/suite/realtime_conversation.rs +++ b/codex-rs/core/tests/suite/realtime_conversation.rs @@ -1,6 +1,10 @@ use anyhow::Context; use anyhow::Result; use chrono::Utc; +use codex_app_server_protocol::ThreadRealtimeItemContent; +use codex_app_server_protocol::ThreadRealtimeSessionOutcome; +use codex_app_server_protocol::ThreadRealtimeTranscriptRole; +use codex_app_server_protocol::ThreadTimelineEntry; use codex_config::config_toml::RealtimeWsMode; use codex_config::config_toml::RealtimeWsVersion; use codex_core::StartThreadOptions; @@ -33,8 +37,10 @@ use codex_protocol::protocol::RealtimeOutputModality; use codex_protocol::protocol::RealtimeTranscriptEntry; use codex_protocol::protocol::RealtimeVoice; use codex_protocol::protocol::SessionSource; +use codex_protocol::protocol::ThreadHistoryMode; use codex_protocol::protocol::ThreadSource; use codex_protocol::user_input::UserInput; +use codex_thread_store::ListTimelineParams; use core_test_support::responses; use core_test_support::responses::WebSocketConnectionConfig; use core_test_support::responses::start_mock_server; @@ -437,10 +443,148 @@ async fn conversation_start_audio_text_close_round_trip() -> Result<()> { Some("requested" | "transport_closed") )); + test.codex.ensure_rollout_materialized().await; + let history = test.codex.load_history(/*include_archived*/ false).await?; + assert!( + !history + .items + .iter() + .any(|item| matches!(item, RolloutItem::RealtimeItem(_))) + ); + server.shutdown().await; Ok(()) } +// No host consumes realtime events in this test. The injected store must receive +// canonical history even when the host has no realtime notification adapter. +#[test_case(ThreadRealtimeSessionOutcome::Ended; "transport_close")] +#[test_case(ThreadRealtimeSessionOutcome::Failed; "error_close")] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn conversation_records_history_without_an_event_observer( + outcome: ThreadRealtimeSessionOutcome, +) -> Result<()> { + skip_if_no_network!(Ok(())); + let api_server = start_mock_server().await; + let mut events = vec![ + json!({ "type": "session.updated", "session": { "id": "voice-1" } }), + json!({ "type": "response.output_text.delta", "delta": "assistant first" }), + json!({ "type": "conversation.item.input_audio_transcription.delta", "delta": "user second" }), + ]; + if outcome == ThreadRealtimeSessionOutcome::Failed { + events.push(json!({ "type": "error", "error": { "message": "fixture failure" } })); + } + let realtime_server = start_websocket_server(vec![vec![events.clone()], vec![events]]).await; + let mut builder = test_codex() + .with_history_mode(ThreadHistoryMode::Paginated) + .with_config({ + let url = realtime_server.uri().to_string(); + move |config| { + config.experimental_realtime_ws_base_url = Some(url); + config.realtime.version = RealtimeWsVersion::V2; + } + }); + let test = builder.build_with_auto_env(&api_server).await?; + test.codex.ensure_rollout_materialized().await; + let mut expected = Vec::new(); + for session_count in 1..=2 { + test.codex + .submit(Op::RealtimeConversationStart(ConversationStartParams { + client_managed_handoffs: false, + delegation_ack_filler: None, + flush_transcript_tail_on_session_end: false, + codex_responses_as_items: false, + codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, + codex_response_handoff_channel_prefixes: None, + model: None, + output_modality: RealtimeOutputModality::Audio, + include_startup_context: false, + initial_items: Vec::new(), + realtime_start_instructions: None, + realtime_end_instructions: None, + prompt: Some(Some("fixture".to_string())), + realtime_session_id: Some("voice-1".to_string()), + transport: None, + version: None, + voice: None, + })) + .await?; + let items = timeout(Duration::from_secs(10), async { + loop { + test.thread_store + .flush_thread(test.session_configured.thread_id) + .await?; + let history = test + .thread_store + .list_timeline(ListTimelineParams { + thread_id: test.session_configured.thread_id, + cursor: None, + page_size: 100, + }) + .await?; + let items = history + .items + .into_iter() + .filter_map(|item| match item { + ThreadTimelineEntry::Realtime { item, .. } => Some(item), + _ => None, + }) + .collect::>(); + if items + .iter() + .filter(|item| { + matches!( + item.content, + ThreadRealtimeItemContent::RealtimeSessionClosed { .. } + ) + }) + .count() + == session_count + { + break Ok::<_, anyhow::Error>(items); + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .context("Core should persist history without a host observer")??; + expected.extend([ + ThreadRealtimeItemContent::RealtimeSessionStarted, + ThreadRealtimeItemContent::TranscriptSegment { + role: ThreadRealtimeTranscriptRole::Assistant, + text: "assistant first".to_string(), + }, + ThreadRealtimeItemContent::TranscriptSegment { + role: ThreadRealtimeTranscriptRole::User, + text: "user second".to_string(), + }, + ThreadRealtimeItemContent::RealtimeSessionClosed { outcome }, + ]); + assert_eq!( + items + .iter() + .map(|item| item.content.clone()) + .collect::>(), + expected + ); + let ids = items + .iter() + .map(|item| item.id.as_str()) + .collect::>(); + assert_eq!(ids.len(), items.len()); + assert!( + items + .iter() + .all(|item| item.realtime_session_id == "voice-1") + ); + } + test.codex.submit(Op::Shutdown).await?; + realtime_server.shutdown().await; + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn conversation_start_defaults_to_v2_and_gpt_realtime_1_5() -> Result<()> { skip_if_no_network!(Ok(())); @@ -3530,8 +3674,10 @@ async fn conversation_startup_context_is_truncated_and_sent_once_per_start() -> Ok(()) } +#[test_case(false; "durable")] +#[test_case(true; "ephemeral")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { +async fn conversation_user_text_turn_is_not_sent_to_realtime(ephemeral: bool) -> Result<()> { skip_if_no_network!(Ok(())); let api_server = start_mock_server().await; @@ -3545,22 +3691,32 @@ async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { .await; let realtime_server = start_websocket_server(vec![vec![ - vec![json!({ - "type": "session.updated", - "session": { "id": "sess_user_text", "instructions": "backend prompt" } - })], + vec![ + json!({ + "type": "session.updated", + "session": { "id": "sess_user_text", "instructions": "backend prompt" } + }), + json!({ + "type": "response.output_text.delta", + "delta": "spoken before typed input" + }), + ], vec![], ]]) .await; - let mut builder = test_codex().with_config({ - let realtime_base_url = realtime_server.uri().to_string(); - move |config| { - config.experimental_realtime_ws_base_url = Some(realtime_base_url); - config.experimental_realtime_ws_startup_context = Some(String::new()); - } - }); - let test = builder.build(&api_server).await?; + let mut builder = test_codex() + .with_history_mode(ThreadHistoryMode::Paginated) + .with_config({ + let realtime_base_url = realtime_server.uri().to_string(); + move |config| { + config.ephemeral = ephemeral; + config.realtime.version = RealtimeWsVersion::V2; + config.experimental_realtime_ws_base_url = Some(realtime_base_url); + config.experimental_realtime_ws_startup_context = Some(String::new()); + } + }); + let test = builder.build_with_auto_env(&api_server).await?; test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { @@ -3599,6 +3755,20 @@ async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { .await; assert_eq!(session_updated, "sess_user_text"); + wait_for_event(&test.codex, |event| { + if let EventMsg::RealtimeConversationRealtime(event) = event { + match event.payload { + RealtimeEvent::HistoryItemStarted(_) + | RealtimeEvent::HistoryTranscriptDelta { .. } + | RealtimeEvent::HistoryItemCompleted(_) => assert!(!ephemeral), + RealtimeEvent::OutputTranscriptDelta(_) => return true, + _ => {} + } + } + false + }) + .await; + let user_text = "typed follow-up for realtime"; test.codex .start_or_steer_turn(TurnInputRequest::user_input(vec![UserInput::Text { @@ -3617,6 +3787,30 @@ async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { let model_user_texts = response_mock.single_request().message_input_texts("user"); assert!(model_user_texts.iter().any(|text| text == user_text)); + if !ephemeral { + test.thread_store + .flush_thread(test.session_configured.thread_id) + .await?; + let timeline = test + .thread_store + .list_timeline(ListTimelineParams { + thread_id: test.session_configured.thread_id, + cursor: None, + page_size: 100, + }) + .await?; + assert_eq!( + timeline.items.iter().filter_map(|entry| match entry { + ThreadTimelineEntry::Realtime { item, .. } + if matches!(&item.content, ThreadRealtimeItemContent::TranscriptSegment { text, .. } if text == "spoken before typed input") => Some("speech"), + ThreadTimelineEntry::Item { item, .. } + if matches!(item.as_ref(), codex_app_server_protocol::ThreadItem::UserMessage { .. }) => Some("typed input"), + _ => None, + }).collect::>(), + vec!["speech", "typed input"] + ); + } + let realtime_connections = realtime_server.connections(); assert_eq!(realtime_connections.len(), 1); assert_eq!(realtime_connections[0].len(), 1); @@ -4246,7 +4440,7 @@ fn message_input_texts(body: &Value, role: &str) -> Vec { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn inbound_handoff_request_starts_turn() -> Result<()> { +async fn inbound_handoff_request_starts_turn_and_promotes_its_artifact() -> Result<()> { skip_if_no_network!(Ok(())); let api_server = start_mock_server().await; @@ -4254,7 +4448,7 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> { &api_server, responses::sse(vec![ responses::ev_response_created("resp-1"), - responses::ev_assistant_message("msg-1", "ok"), + responses::ev_assistant_message("msg-1", "::codex-realtime-inline{}\nShared artifact"), responses::ev_completed("resp-1"), ]), ) @@ -4278,14 +4472,16 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> { ]]]) .await; - let mut builder = test_codex().with_config({ - let realtime_base_url = realtime_server.uri().to_string(); - move |config| { - config.experimental_realtime_ws_base_url = Some(realtime_base_url); - config.realtime.version = RealtimeWsVersion::V1; - } - }); - let test = builder.build(&api_server).await?; + let mut builder = test_codex() + .with_history_mode(ThreadHistoryMode::Paginated) + .with_config({ + let realtime_base_url = realtime_server.uri().to_string(); + move |config| { + config.experimental_realtime_ws_base_url = Some(realtime_base_url); + config.realtime.version = RealtimeWsVersion::V1; + } + }); + let test = builder.build_with_auto_env(&api_server).await?; test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { @@ -4349,6 +4545,42 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> { }) .await; + test.thread_store + .flush_thread(test.session_configured.thread_id) + .await?; + let timeline = test + .thread_store + .list_timeline(ListTimelineParams { + thread_id: test.session_configured.thread_id, + cursor: None, + page_size: 100, + }) + .await?; + let promotions = timeline + .items + .into_iter() + .filter_map(|entry| match entry { + ThreadTimelineEntry::Realtime { item, .. } + if matches!( + item.content, + ThreadRealtimeItemContent::BemItemPromoted { .. } + ) => + { + Some(item.content) + } + _ => None, + }) + .collect::>(); + assert_eq!( + promotions, + vec![ThreadRealtimeItemContent::BemItemPromoted { + turn_id, + item_id: "msg-1".to_string(), + presentation: + codex_app_server_protocol::ThreadRealtimeBemItemPresentation::InlineMarkdown, + }] + ); + let request = response_mock.single_request(); let turn_metadata: Value = serde_json::from_str( request diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index 06ba9705f7b7..d4d236d963c2 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -448,6 +448,13 @@ pub enum RealtimeEvent { ConversationItemDone { item_id: String, }, + /// Canonical display history produced by Core, separate from provider events. + HistoryItemStarted(crate::realtime::RealtimeItem), + HistoryTranscriptDelta { + item_id: String, + delta: String, + }, + HistoryItemCompleted(crate::realtime::RealtimeItem), HandoffRequested(RealtimeHandoffRequested), NoopRequested(RealtimeNoopRequested), Error(String),