From 0e8bc4f8a14acd0ef1f395f95199ddc820961911 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 21 Jul 2026 17:45:12 +0200 Subject: [PATCH 01/28] Core: Request history on channel join and create batched history event --- packages/core/core-shared/src/actor.rs | 397 ++++++++++++------ packages/core/core-shared/src/send_command.rs | 69 +++ packages/core/core-shared/src/state.rs | 90 +++- packages/core/core-wasm/src/lib.rs | 35 +- 4 files changed, 449 insertions(+), 142 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index ea29a6e..f89837b 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -6,7 +6,7 @@ use crate::{ SendCommand, state::{ Channel, ChannelRole, ChannelUser, Message, MessageMetadata, MessageReference, MessageType, - OrbitError, React, Server, ServerEvent, SignedIn, TextMessage, User, + OrbitError, Server, ServerEvent, SignedIn, Tags, TextMessage, User, }, }; use anyhow::{Context, anyhow}; @@ -19,11 +19,8 @@ use futures::{ }, stream::FusedStream, }; -use irc_proto::{ - BatchSubCommand, CapSubCommand, Command::*, Message as IrcMessage, Response, message::Tag, -}; +use irc_proto::{BatchSubCommand, CapSubCommand, Command::*, Message as IrcMessage, Response}; use ordermap::OrderMap; -use time::{OffsetDateTime, format_description::well_known::Iso8601}; use tracing::{debug, error, warn}; #[derive(Debug, Default)] @@ -115,24 +112,12 @@ pub trait IrcConnection: fmt::Debug { fn address(&self) -> &str; } -pub struct IrcActor { - cmd_rx: mpsc::UnboundedReceiver, - incoming: C::Incoming, - outgoing: C::Outgoing, - state: Server, - response_channels: ResponseChannels, - event_handlers: Vec>, - error_handlers: Vec>, - disconnect_handlers: Vec>, - - current_batch: Option, - sasl_state: SaslState, -} - #[derive(Debug)] struct Batch { id: String, typ: BatchSubCommand, + channel: String, + messages: Vec, } impl Batch { @@ -153,6 +138,20 @@ enum SaslState { }, } +pub struct IrcActor { + cmd_rx: mpsc::UnboundedReceiver, + incoming: C::Incoming, + outgoing: C::Outgoing, + state: Server, + response_channels: ResponseChannels, + event_handlers: Vec>, + error_handlers: Vec>, + disconnect_handlers: Vec>, + + current_batch: Option, + sasl_state: SaslState, +} + impl IrcActor { #[tracing::instrument] pub async fn start( @@ -236,7 +235,32 @@ impl IrcActor { JOIN(ref channel_name, _, _) => { let source = message.source_nickname().unwrap(); - // FIXME: handle other cases + let mut tags = Tags::default(); + if let Some(ref t) = message.tags { + tags = Tags::parse(t); + } + + if self.current_batch.as_ref().map(|b| b.is_chathistory()) == Some(true) { + assert!(tags.server_time.is_some()); + } + + let state_message = Message { + text: None, + metadata: MessageMetadata { + msgid: tags.msgid_with_fallback(&[source]), + server_time: tags.server_time_with_fallback() as f64, + message_type: MessageType::Join, + user: source.to_string(), + }, + }; + + if self + .push_batch(channel_name.clone(), state_message.clone()) + .await + { + return Ok(()); + } + if source == self.state.me.as_ref().unwrap().nickname { let channel = Channel::new(channel_name.clone()); self.state @@ -251,88 +275,168 @@ impl IrcActor { .map_err(|e| anyhow!("Failed to reply to JOIN command {e:?}"))?; self.on_event(ServerEvent::Joined(channel)).await?; + + if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) + && self.state.capabilities.history.enabled + { + self.history_latest(channel_name.clone(), None, 50) + .await + .context("Failed to request latest history")?; + } + } else { + if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { + let channel = self.channel_mut(channel_name.clone()).await; + + dbg!(&message); + channel.users.push(ChannelUser { + nickname: source.to_string(), + role: ChannelRole::None, + }); + } } + + self.on_event(ServerEvent::Privmsg { + channel: channel_name.to_string(), + message: state_message, + }) + .await?; } - PRIVMSG(ref target, ref text) => { - let mut msgid = None; - let mut server_time = None; - let mut username = None; - let mut relayed_by = None; - let mut reply = None; - if let Some(ref tags) = message.tags { - for Tag(key, value) in tags { - match key.as_str() { - "msgid" => msgid = value.clone(), - "account" => username = value.clone(), - "draft/relaymsg" => relayed_by = value.clone(), - "+draft/reply" | "+reply" => reply = value.clone(), - "time" => { - server_time = value - .as_ref() - .and_then(|v| OffsetDateTime::parse(v, &Iso8601::DEFAULT).ok()) - } - _ => { - warn!("unhandled tag: {key:?}: {value:?}"); - } - } - } + PART(ref channel_name, ref comment) => { + let mut tags = Tags::default(); + if let Some(ref t) = message.tags { + tags = Tags::parse(t); } + let source = message.source_nickname().unwrap(); - let nickname = message.source_nickname().unwrap(); + let state_message = Message { + text: None, + metadata: MessageMetadata { + msgid: tags.msgid_with_fallback(&[source]), + server_time: tags.server_time_with_fallback() as f64, + message_type: MessageType::Part, + user: source.to_string(), + }, + }; - if let Some(username) = username { - let user = self - .state - .users - .entry(nickname.to_string()) - .or_insert_with(|| User::new(nickname.to_string())); - user.username = Some(username); + if self + .push_batch(channel_name.clone(), state_message.clone()) + .await + { + return Ok(()); + } + + self.on_event(ServerEvent::Privmsg { + channel: channel_name.to_string(), + message: state_message, + }) + .await?; + + if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { + let channel = self.channel_mut(channel_name.clone()).await; + + channel.users.retain(|u| u.nickname != source); + } + } + QUIT(ref comment) => { + let mut tags = Tags::default(); + if let Some(ref t) = message.tags { + tags = Tags::parse(t); + } + let source = message.source_nickname().unwrap(); + + let state_message = Message { + text: None, + metadata: MessageMetadata { + msgid: tags.msgid_with_fallback(&[source]), + server_time: tags.server_time_with_fallback() as f64, + message_type: MessageType::Part, + user: source.to_string(), + }, + }; + + if let Some(batch) = self.current_batch.as_mut() + && batch.is_chathistory() + { + batch.messages.push(state_message); + return Ok(()); } - let server_time = server_time - .unwrap_or_else(OffsetDateTime::now_utc) - .unix_timestamp(); - let msgid = msgid.unwrap_or_else(|| { - let mut hasher = blake3::Hasher::new(); - hasher.update(&server_time.to_ne_bytes()); - hasher.update(target.as_bytes()); - hasher.update(text.as_bytes()); + self.on_event(ServerEvent::Privmsg { + channel: String::new(), + message: state_message, + }) + .await?; + + if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { + self.state.users.remove(source); + for channel in self.state.channels.values_mut() { + channel.users.retain(|u| u.nickname != source); + } + } + } + PRIVMSG(ref target, ref text) => { + let mut tags = Tags::default(); + if let Some(ref t) = message.tags { + tags = Tags::parse(t); + } + assert_eq!( + self.current_batch.as_ref().map(|s| s.id.as_str()), + tags.batch.as_deref(), + ); + + let source = message.source_nickname().unwrap(); - hasher.finalize().to_string() - }); + let msgid = tags.msgid_with_fallback(&[source, target, text]); - let reply = reply - .and_then(|r| { + let reply = tags + .reply + .as_ref() + .map(|r| { self.state .channels .get(target) - .and_then(|c| c.messages.get(&r)) + .and_then(|c| c.messages.get(r)) }) - .and_then(|m| { - Some(MessageReference { - text: m.text.clone().map(|t| t.content)?, - username: m.metadata.user.clone(), - }) + .map(|m| MessageReference { + text: m + .map(|m| m.text.clone().map(|t| t.content)) + .unwrap_or_else(|| Some(String::from("Unknown"))), + + username: m + .map(|m| m.metadata.user.clone()) + .unwrap_or_else(|| String::from("Unknown")), }); let state_message = Message { + metadata: MessageMetadata { + msgid: msgid.clone(), + server_time: tags.server_time_with_fallback() as f64, + message_type: MessageType::Privmsg, + user: source.to_string(), + }, text: Some(TextMessage { content: text.clone(), reactions: OrderMap::new(), reply, redacted: false, edited: false, - relayed_by, + relayed_by: tags.relayed_by, }), - metadata: MessageMetadata { - msgid: msgid.clone(), - server_time: server_time as f64, - message_type: MessageType::Privmsg, - user: nickname.to_string(), - }, }; - if nickname == self.state.me.as_ref().unwrap().nickname + if self.push_batch(target.clone(), state_message.clone()).await { + return Ok(()); + } + + if let Some(username) = tags.username.clone() { + let user = self.user_mut(source.to_string()).await; + user.username = Some(username); + } + + let channel = self.channel_mut(target.clone()).await; + channel.messages.insert(msgid, state_message.clone()); + + if source == self.state.me.as_ref().unwrap().nickname && let Err(e) = self .response_channels .reply( @@ -347,30 +451,24 @@ impl IrcActor { error!("Failed to reply to PRIVMSG command {e:?}"); } - let channel = self - .state - .channels - .entry(target.clone()) - .or_insert_with(|| Channel::new(target.clone())); - channel.messages.insert(msgid, state_message.clone()); - - if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { - self.on_event(ServerEvent::Privmsg { - channel: target.clone(), - message: state_message, - }) - .await?; - } + self.on_event(ServerEvent::Privmsg { + channel: target.clone(), + message: state_message, + }) + .await?; } BATCH(reference, typ, param) => { if let Some(id) = reference.strip_prefix('+') { self.current_batch = Some(Batch { id: id.to_string(), typ: typ.clone().unwrap(), + channel: String::new(), + messages: Vec::new(), }); match typ { - Some(BatchSubCommand::CUSTOM(c)) if &c == "METADATA" => (), + Some(BatchSubCommand::CUSTOM(c)) + if &c == "METADATA" || &c == "CHATHISTORY" => {} _ => warn!(?typ, ?param, "unhandled BATCH type"), } } else { @@ -378,39 +476,62 @@ impl IrcActor { self.current_batch.as_ref().map(|s| s.id.as_str()), Some(&reference[1..]) ); + // FIXME: this will panic if a batch only contains QUITs + assert_eq!( + self.current_batch.as_ref().map(|b| b.channel.is_empty()), + Some(false) + ); + + if let Some(batch) = self.current_batch.take() + && !batch.messages.is_empty() + { + for message in &batch.messages { + let channel = self.channel_mut(batch.channel.clone()).await; + channel + .messages + .insert(message.metadata.msgid.clone(), message.clone()); + } + + self.on_event(ServerEvent::History { + channel: batch.channel, + messages: batch.messages, + }) + .await?; + } self.current_batch = None; } } + ChannelMODE(ref channel_name, ref mode) => { + let target = message.response_target().unwrap(); + let source = message.source_nickname().unwrap(); + + if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { + dbg!(target, source, channel_name, mode); + } + } + TOPIC(ref channel_name, ref text) => { + let target = message.response_target().unwrap(); + let source = message.source_nickname().unwrap(); + + if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { + dbg!(target, source, channel_name, text); + } + } Raw(ref cmd, ref mut target) if cmd == "TAGMSG" => { let target = target.remove(0); - let mut react = None; - let mut unreact = None; - let mut reply = None; - if let Some(ref tags) = message.tags { - for Tag(key, value) in tags { - match key.as_str() { - "+draft/reply" | "+reply" => reply = value.clone(), - "+draft/react" => react = value.clone(), - "+draft/unreact" => unreact = value.clone(), - _ => { - warn!("unhandled tag: {key:?}: {value:?}"); - } - } - } + let mut tags = Tags::default(); + if let Some(ref t) = message.tags { + tags = Tags::parse(t); } - let channel = self - .state - .channels - .entry(target.clone()) - .or_insert_with(|| Channel::new(target.clone())); - - let is_unreact = unreact.is_some(); - if let Some(react) = react.or(unreact) - && let Some(reply) = reply + let is_unreact = tags.unreact.is_some(); + if let Some(react) = tags.react.or(tags.unreact) + && let Some(reply) = tags.reply { + let channel = self.channel_mut(target.clone()).await; + let nickname = message.source_nickname().unwrap().to_string(); if let Some(message) = channel.messages.get_mut(&reply) { let reactors = message @@ -428,12 +549,12 @@ impl IrcActor { } // TODO: should it be sent if the message wasn't found? - self.on_event(ServerEvent::React(React { + self.on_event(ServerEvent::React { target_message: reply, user: nickname, text: react, is_unreact, - })) + }) .await?; } } @@ -491,6 +612,22 @@ impl IrcActor { Ok(()) } + pub async fn push_batch(&mut self, target: String, state_message: Message) -> bool { + if let Some(batch) = self.current_batch.as_mut() + && batch.is_chathistory() + { + if batch.channel.is_empty() { + batch.channel = target.clone(); + } else { + assert_eq!(batch.channel, target) + } + batch.messages.push(state_message); + + return true; + } + return false; + } + #[tracing::instrument(err, skip(self))] pub async fn handle_caps( &mut self, @@ -598,11 +735,7 @@ impl IrcActor { Response::RPL_TOPIC => { let channel_name = params[1].to_string(); let topic = params[2].to_string(); - let channel = self - .state - .channels - .entry(channel_name.clone()) - .or_insert_with(|| Channel::new(channel_name)); + let channel = self.channel_mut(channel_name).await; channel.metadata.topic = Some(topic); @@ -640,11 +773,7 @@ impl IrcActor { .or_insert_with(|| User::new(user)); } - let channel = self - .state - .channels - .entry(channel_name.clone()) - .or_insert_with(|| Channel::new(channel_name)); + let channel = self.channel_mut(channel_name).await; channel.users = channel_users; } @@ -843,6 +972,20 @@ impl IrcActor { Ok(()) } + + async fn channel_mut(&mut self, name: String) -> &mut Channel { + self.state + .channels + .entry(name.clone()) + .or_insert_with(|| Channel::new(name)) + } + + async fn user_mut(&mut self, nickname: String) -> &mut User { + self.state + .users + .entry(nickname.clone()) + .or_insert_with(|| User::new(nickname)) + } } impl SendCommand for IrcActor { diff --git a/packages/core/core-shared/src/send_command.rs b/packages/core/core-shared/src/send_command.rs index 373a05a..9545f37 100644 --- a/packages/core/core-shared/src/send_command.rs +++ b/packages/core/core-shared/src/send_command.rs @@ -1,3 +1,5 @@ +use core::fmt; + use futures::{ SinkExt, channel::mpsc::{self, UnboundedSender}, @@ -108,6 +110,51 @@ pub trait SendCommand { ) -> impl std::future::Future> { async { self.command(WHOIS(server, user)).await } } + + fn history( + &mut self, + subcommand: ChatHistorySubCommand, + mut args: Vec, + ) -> impl std::future::Future> { + async move { + args.insert(0, subcommand.to_string()); + self.command(Raw("CHATHISTORY".to_string(), args)).await + } + } + + fn history_before( + &mut self, + target: String, + before: String, + limit: i32, + ) -> impl std::future::Future> { + async move { + self.history( + ChatHistorySubCommand::Before, + vec![target, before, limit.to_string()], + ) + .await + } + } + + fn history_latest( + &mut self, + target: String, + since: Option, + limit: i32, + ) -> impl std::future::Future> { + async move { + self.history( + ChatHistorySubCommand::Latest, + vec![ + target, + since.unwrap_or_else(|| String::from("*")), + limit.to_string(), + ], + ) + .await + } + } } impl SendCommand for UnboundedSender { @@ -118,3 +165,25 @@ impl SendCommand for UnboundedSender { Ok(()) } } + +pub enum ChatHistorySubCommand { + Before, + After, + Latest, + Around, + Between, + Targets, +} + +impl fmt::Display for ChatHistorySubCommand { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Before => write!(f, "BEFORE"), + Self::After => write!(f, "AFTER"), + Self::Latest => write!(f, "LATEST"), + Self::Around => write!(f, "AROUND"), + Self::Between => write!(f, "BETWEEN"), + Self::Targets => write!(f, "TARGETS"), + } + } +} diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index 8c67dbc..cd67e8b 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -3,9 +3,12 @@ use std::str::FromStr; #[cfg(feature = "web")] use crate::dbg; +use irc_proto::message::Tag; use ordermap::OrderMap; use thiserror::Error; -use tracing::error; +use time::OffsetDateTime; +use time::format_description::well_known::Iso8601; +use tracing::{error, warn}; #[cfg(feature = "web")] use tsify::Tsify; #[cfg(feature = "web")] @@ -547,8 +550,8 @@ impl Eq for MessageMetadata {} #[cfg_attr(feature = "web", derive(Tsify))] #[cfg_attr(feature = "web", wasm_bindgen(getter_with_clone, inspectable))] pub struct MessageReference { + pub text: Option, pub username: String, - pub text: String, } #[derive(Debug, Clone, PartialEq, Eq)] @@ -564,17 +567,16 @@ pub enum ServerEvent { channel: String, message: Message, }, - React(React), -} - -#[derive(Debug, Clone, PartialEq, Eq)] -#[cfg_attr(feature = "web", derive(Tsify))] -#[cfg_attr(feature = "web", wasm_bindgen(getter_with_clone, inspectable))] -pub struct React { - pub target_message: String, - pub user: String, - pub text: String, - pub is_unreact: bool, + React { + target_message: String, + user: String, + text: String, + is_unreact: bool, + }, + History { + channel: String, + messages: Vec, + }, } #[derive(Debug, Clone)] @@ -610,3 +612,65 @@ impl From for OrbitError { Self::Unknown(error.to_string()) } } + +#[derive(Debug, Default)] +pub struct Tags { + pub msgid: Option, + pub server_time: Option, + pub username: Option, + pub relayed_by: Option, + pub reply: Option, + pub react: Option, + pub unreact: Option, + pub batch: Option, + pub typing: Option, + pub bot: Option, +} + +impl Tags { + pub fn parse(tags: &Vec) -> Self { + let mut out = Tags::default(); + + for Tag(key, value) in tags { + match key.as_str() { + "msgid" => out.msgid = value.clone(), + "account" => out.username = value.clone(), + "draft/relaymsg" => out.relayed_by = value.clone(), + "batch" => out.batch = value.clone(), + "bot" => out.bot = value.clone(), + "time" => { + out.server_time = value + .as_ref() + .and_then(|v| OffsetDateTime::parse(v, &Iso8601::DEFAULT).ok()) + } + "+draft/reply" | "+reply" => out.reply = value.clone(), + "+draft/react" => out.react = value.clone(), + "+draft/unreact" => out.unreact = value.clone(), + "+typing" => out.typing = value.clone(), + _ => { + warn!("unhandled tag: {key:?}: {value:?}"); + } + } + } + + out + } + + pub fn server_time_with_fallback(&self) -> i64 { + self.server_time + .unwrap_or_else(OffsetDateTime::now_utc) + .unix_timestamp() + } + + pub fn msgid_with_fallback(&self, hash_extras: &[&str]) -> String { + self.msgid.clone().unwrap_or_else(|| { + let mut hasher = blake3::Hasher::new(); + hasher.update(&self.server_time_with_fallback().to_ne_bytes()); + for extra in hash_extras { + hasher.update(extra.as_bytes()); + } + + hasher.finalize().to_string() + }) + } +} diff --git a/packages/core/core-wasm/src/lib.rs b/packages/core/core-wasm/src/lib.rs index f79f1c6..ede4014 100644 --- a/packages/core/core-wasm/src/lib.rs +++ b/packages/core/core-wasm/src/lib.rs @@ -5,7 +5,7 @@ use core_shared::{ SendCommand, actor::{self, ActorCommand, ActorMessage, CommandResponse, IrcActor}, state::{ - self, Capabilities, ChannelMetadata, ChannelUser, MessageMetadata, MessageReference, React, + self, Capabilities, ChannelMetadata, ChannelUser, MessageMetadata, MessageReference, ServerMetadata, SignedIn, User, }, }; @@ -475,6 +475,7 @@ pub enum ServerEvent { UserList(UserList), Privmsg(ChannelMessage), React(React), + History(History), } impl From for ServerEvent { @@ -490,7 +491,21 @@ impl From for ServerEvent { channel, message: message.into(), }), - state::ServerEvent::React(r) => Self::React(r), + state::ServerEvent::React { + target_message, + user, + text, + is_unreact, + } => Self::React(React { + target_message, + user, + text, + is_unreact, + }), + state::ServerEvent::History { channel, messages } => Self::History(History { + channel, + messages: messages.into_iter().map(|m| m.into()).collect(), + }), } } } @@ -595,3 +610,19 @@ impl From for OrbitError { } } } + +#[derive(Debug, Clone, PartialEq, Eq, Tsify)] +#[wasm_bindgen(getter_with_clone, inspectable)] +pub struct React { + pub target_message: String, + pub user: String, + pub text: String, + pub is_unreact: bool, +} + +#[derive(Debug, Clone, Tsify)] +#[wasm_bindgen(getter_with_clone, inspectable)] +pub struct History { + pub channel: String, + pub messages: Vec, +} From 2078a99ab94bcd6c8809af1099226c8d04a185ed Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 21 Jul 2026 18:30:05 +0200 Subject: [PATCH 02/28] Add untested history_before command --- packages/core/core-shared/src/actor.rs | 36 +++++++++++++++++++---- packages/core/core-shared/src/state.rs | 11 ++++--- packages/core/core-wasm/src/lib.rs | 40 +++++++++++++++++++++++--- 3 files changed, 73 insertions(+), 14 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index f89837b..9fde6b2 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -5,8 +5,8 @@ use crate::dbg; use crate::{ SendCommand, state::{ - Channel, ChannelRole, ChannelUser, Message, MessageMetadata, MessageReference, MessageType, - OrbitError, Server, ServerEvent, SignedIn, Tags, TextMessage, User, + Channel, ChannelRole, ChannelUser, History, Message, MessageMetadata, MessageReference, + MessageType, OrbitError, Server, ServerEvent, SignedIn, Tags, TextMessage, User, }, }; use anyhow::{Context, anyhow}; @@ -32,6 +32,7 @@ pub enum CommandKey { SignIn, Join(String), Privmsg { target: String, text: String }, + History { target: String }, } #[derive(Debug)] @@ -41,6 +42,7 @@ pub enum CommandResponse { SignIn(Result), Join(String), Privmsg(Box), + History(History), } impl ResponseChannels { @@ -102,6 +104,10 @@ pub enum ActorCommand { AddDisconectHandler { handler: UnboundedSender, }, + RequestHistory { + channel: String, + before_msgid: String, + }, } pub trait IrcConnection: fmt::Debug { @@ -492,11 +498,22 @@ impl IrcActor { .insert(message.metadata.msgid.clone(), message.clone()); } - self.on_event(ServerEvent::History { - channel: batch.channel, + let history = History { + channel: batch.channel.clone(), messages: batch.messages, - }) - .await?; + }; + + self.response_channels + .reply( + &CommandKey::History { + target: batch.channel, + }, + CommandResponse::History(history.clone()), + ) + .await + .unwrap(); + + self.on_event(ServerEvent::History(history)).await?; } self.current_batch = None; @@ -862,6 +879,13 @@ impl IrcActor { ActorCommand::AddDisconectHandler { handler } => { self.disconnect_handlers.push(handler); } + ActorCommand::RequestHistory { + channel, + before_msgid, + } => self + .history_before(channel, before_msgid, 50) + .await + .context("Failed to send history before")?, } Ok(()) diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index cd67e8b..a9597d6 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -573,10 +573,13 @@ pub enum ServerEvent { text: String, is_unreact: bool, }, - History { - channel: String, - messages: Vec, - }, + History(History), +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct History { + pub channel: String, + pub messages: Vec, } #[derive(Debug, Clone)] diff --git a/packages/core/core-wasm/src/lib.rs b/packages/core/core-wasm/src/lib.rs index ede4014..276d521 100644 --- a/packages/core/core-wasm/src/lib.rs +++ b/packages/core/core-wasm/src/lib.rs @@ -315,6 +315,32 @@ impl IrcConnection { address: self.address.clone(), }) } + + #[wasm_bindgen] + pub async fn history_before( + &mut self, + channel: String, + before_msgid: String, + ) -> Result { + let (tx, rx) = oneshot::channel(); + self.address + .send(ActorMessage { + command: ActorCommand::RequestHistory { + channel, + before_msgid, + }, + reply_tx: Some(tx), + }) + .await + .context("Failed to send ActorMessage")?; + + let resp = rx.await.context("Failed to await ActorMessage")?; + let CommandResponse::History(history) = resp else { + unreachable!("expected history, got: {:?}", resp); + }; + + Ok(history.into()) + } } #[wasm_bindgen] @@ -502,10 +528,7 @@ impl From for ServerEvent { text, is_unreact, }), - state::ServerEvent::History { channel, messages } => Self::History(History { - channel, - messages: messages.into_iter().map(|m| m.into()).collect(), - }), + state::ServerEvent::History(history) => Self::History(history.into()), } } } @@ -626,3 +649,12 @@ pub struct History { pub channel: String, pub messages: Vec, } + +impl From for History { + fn from(history: state::History) -> Self { + Self { + channel: history.channel, + messages: history.messages.into_iter().map(Into::into).collect(), + } + } +} From f6ae112fd42c7b2d1387e12c16e4567c165ef9f7 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 21 Jul 2026 18:42:58 +0200 Subject: [PATCH 03/28] Fix incomplete assert --- packages/core/core-shared/src/actor.rs | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 9fde6b2..72d4a70 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -484,7 +484,9 @@ impl IrcActor { ); // FIXME: this will panic if a batch only contains QUITs assert_eq!( - self.current_batch.as_ref().map(|b| b.channel.is_empty()), + self.current_batch + .as_ref() + .map(|b| b.is_chathistory() && b.channel.is_empty()), Some(false) ); @@ -642,7 +644,8 @@ impl IrcActor { return true; } - return false; + + false } #[tracing::instrument(err, skip(self))] From 1bee42fd33af3cb945dd49d92a9e89b181f67bd2 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 21 Jul 2026 23:58:55 +0200 Subject: [PATCH 04/28] Time out response channels after ~1s --- packages/core/Cargo.lock | 17 +++++- packages/core/core-shared/Cargo.toml | 13 ++-- packages/core/core-shared/src/actor.rs | 85 +++++++++++++++----------- packages/core/core-shared/src/state.rs | 13 ++++ packages/core/core-wasm/Cargo.toml | 2 +- 5 files changed, 89 insertions(+), 41 deletions(-) diff --git a/packages/core/Cargo.lock b/packages/core/Cargo.lock index a1f5c91..14daebf 100644 --- a/packages/core/Cargo.lock +++ b/packages/core/Cargo.lock @@ -128,7 +128,7 @@ dependencies = [ "futures", "gloo-console 0.4.0", "gloo-net 0.7.0", - "indexed_db_futures", + "gloo-timers 0.4.0", "irc-proto", "ordermap", "thiserror 2.0.18", @@ -137,6 +137,7 @@ dependencies = [ "tracing", "tsify", "wasm-bindgen", + "web-time", ] [[package]] @@ -354,7 +355,7 @@ dependencies = [ "gloo-net 0.3.1", "gloo-render", "gloo-storage", - "gloo-timers", + "gloo-timers 0.2.6", "gloo-utils 0.1.7", "gloo-worker", ] @@ -510,6 +511,18 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "gloo-timers" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "482ce8a491a501da4cd806bd190275363d674f2845005c6ddbd5d3e1dd54495d" +dependencies = [ + "futures-channel", + "futures-core", + "js-sys", + "wasm-bindgen", +] + [[package]] name = "gloo-utils" version = "0.1.7" diff --git a/packages/core/core-shared/Cargo.toml b/packages/core/core-shared/Cargo.toml index ac3cff8..b34f2ba 100644 --- a/packages/core/core-shared/Cargo.toml +++ b/packages/core/core-shared/Cargo.toml @@ -1,4 +1,5 @@ [package] + name = "core-shared" version = "0.1.0" edition = "2024" @@ -8,19 +9,23 @@ anyhow = { version = "1.0.103", features = ["backtrace"] } base64 = "0.22.1" blake3 = "1.8.5" futures = "0.3.32" -gloo-console = { version = "0.4.0", optional = true } -gloo-net = { version = "0.7.0", optional = true } -indexed_db_futures = { version = "0.6.4", features = ["serde"] } irc-proto = "1.1.0" ordermap = "1.2.0" thiserror = "2.0.18" time = { version = "0.3.53", features = ["parsing"] } tracing = "0.1.44" + tsify = { version = "0.5.6", optional = true, features = ["js"] } wasm-bindgen = { version = "0.2.126", optional = true } +web-time = { version = "1.1.0", optional = true } +gloo-console = { version = "0.4.0", optional = true } +gloo-net = { version = "0.7.0", optional = true } +tokio = { version = "1.52.3", features = ["macros", "rt", "time", "sync"], optional = true } +gloo-timers = { version = "0.4.0", features = ["futures"], optional = true } [features] -web = ["dep:gloo-net", "dep:wasm-bindgen", "dep:gloo-console", "dep:tsify"] +web = ["dep:gloo-net", "dep:wasm-bindgen", "dep:gloo-console", "dep:tsify", "dep:web-time", "dep:gloo-timers"] +default = ["dep:tokio"] [dev-dependencies] assert_matches = "1.5.0" diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 72d4a70..160d09c 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -1,4 +1,12 @@ -use std::fmt; +#[cfg(feature = "default")] +use std::pin::pin; +use std::{fmt, time::Duration}; + +use futures::FutureExt; +#[cfg(not(feature = "web"))] +use std::time::Instant; +#[cfg(feature = "web")] +use web_time::Instant; #[cfg(feature = "web")] use crate::dbg; @@ -24,7 +32,7 @@ use ordermap::OrderMap; use tracing::{debug, error, warn}; #[derive(Debug, Default)] -pub struct ResponseChannels(Vec<(CommandKey, oneshot::Sender)>); +pub struct ResponseChannels(Vec<(CommandKey, Instant, oneshot::Sender)>); #[derive(Debug, Clone, PartialEq, Eq)] pub enum CommandKey { @@ -47,17 +55,17 @@ pub enum CommandResponse { impl ResponseChannels { pub fn register(&mut self, key: CommandKey, os_tx: oneshot::Sender) { - self.0.push((key, os_tx)); + self.0.push((key, Instant::now(), os_tx)); } #[tracing::instrument] - pub async fn reply( + pub fn reply( &mut self, key: &CommandKey, response: CommandResponse, ) -> Result<(), CommandResponse> { - if let Some(idx) = self.0.iter().position(|(rk, _)| rk == key) { - let (_, ch) = self.0.remove(idx); + if let Some(idx) = self.0.iter().position(|(rk, _, _)| rk == key) { + let (_, _, ch) = self.0.remove(idx); ch.send(response)?; } else { // warn!("Failed to find response channel"); @@ -65,6 +73,11 @@ impl ResponseChannels { Ok(()) } + + pub fn check_timeouts(&mut self) { + self.0 + .retain(|(_, creation, _)| creation.elapsed() < Duration::from_secs(1)); + } } #[derive(Debug)] @@ -198,6 +211,11 @@ impl IrcActor { #[tracing::instrument(skip(self))] pub async fn run(mut self) { loop { + #[cfg(feature = "web")] + let mut timeout = gloo_timers::future::TimeoutFuture::new(1000).fuse(); + #[cfg(not(feature = "web"))] + let mut timeout = pin!(tokio::time::sleep(Duration::from_secs(1)).fuse()); + futures::select! { msg = self.incoming.next() => { match msg { @@ -219,6 +237,9 @@ impl IrcActor { cmd = self.cmd_rx.select_next_some() => { self.handle_command(cmd).await.unwrap(); } + _ = timeout => { + self.response_channels.check_timeouts(); + } } } } @@ -277,7 +298,6 @@ impl IrcActor { &CommandKey::Join(channel_name.clone()), CommandResponse::Join(channel_name.clone()), ) - .await .map_err(|e| anyhow!("Failed to reply to JOIN command {e:?}"))?; self.on_event(ServerEvent::Joined(channel)).await?; @@ -293,7 +313,6 @@ impl IrcActor { if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { let channel = self.channel_mut(channel_name.clone()).await; - dbg!(&message); channel.users.push(ChannelUser { nickname: source.to_string(), role: ChannelRole::None, @@ -443,16 +462,13 @@ impl IrcActor { channel.messages.insert(msgid, state_message.clone()); if source == self.state.me.as_ref().unwrap().nickname - && let Err(e) = self - .response_channels - .reply( - &CommandKey::Privmsg { - target: target.clone(), - text: text.clone(), - }, - CommandResponse::Privmsg(Box::new(state_message.clone())), - ) - .await + && let Err(e) = self.response_channels.reply( + &CommandKey::Privmsg { + target: target.clone(), + text: text.clone(), + }, + CommandResponse::Privmsg(Box::new(state_message.clone())), + ) { error!("Failed to reply to PRIVMSG command {e:?}"); } @@ -482,17 +498,16 @@ impl IrcActor { self.current_batch.as_ref().map(|s| s.id.as_str()), Some(&reference[1..]) ); + // FIXME: this will panic if a batch only contains QUITs assert_eq!( - self.current_batch - .as_ref() - .map(|b| b.is_chathistory() && b.channel.is_empty()), + self.current_batch.as_ref().map(|b| b.is_chathistory() + && !b.messages.is_empty() + && b.messages.is_empty()), Some(false) ); - if let Some(batch) = self.current_batch.take() - && !batch.messages.is_empty() - { + if let Some(batch) = self.current_batch.take() { for message in &batch.messages { let channel = self.channel_mut(batch.channel.clone()).await; channel @@ -512,7 +527,6 @@ impl IrcActor { }, CommandResponse::History(history.clone()), ) - .await .unwrap(); self.on_event(ServerEvent::History(history)).await?; @@ -683,7 +697,6 @@ impl IrcActor { } self.response_channels .reply(&CommandKey::RequestCaps, CommandResponse::Capabilities) - .await .unwrap(); } _ => { @@ -717,7 +730,6 @@ impl IrcActor { &CommandKey::SignIn, CommandResponse::SignIn(Ok(SignedIn::User)), ) - .await .map_err(|e| anyhow!("Failed to reply to sign in command {e:?}"))?; self.cap_end().await.context("Failed to send CAP END")?; @@ -728,7 +740,6 @@ impl IrcActor { &CommandKey::SignIn, CommandResponse::SignIn(Ok(SignedIn::Guest)), ) - .await .map_err(|e| anyhow!("Failed to reply to sign in command {e:?}"))?; } Response::RPL_LOGGEDIN => { @@ -740,7 +751,6 @@ impl IrcActor { &CommandKey::SignIn, CommandResponse::SignIn(Err(OrbitError::SaslFailed(params[1].to_string()))), ) - .await .map_err(|e| anyhow!("Failed to reply to sign in command {e:?}"))?; } Response::ERR_NICKNAMEINUSE => { @@ -749,7 +759,6 @@ impl IrcActor { &CommandKey::SignIn, CommandResponse::SignIn(Err(OrbitError::NickTaken)), ) - .await .map_err(|e| anyhow!("Failed to reply to sign in command {e:?}"))?; } Response::RPL_TOPIC => { @@ -885,10 +894,18 @@ impl IrcActor { ActorCommand::RequestHistory { channel, before_msgid, - } => self - .history_before(channel, before_msgid, 50) - .await - .context("Failed to send history before")?, + } => { + self.response_channels.register( + CommandKey::History { + target: channel.clone(), + }, + cmd.reply_tx.unwrap(), + ); + + self.history_before(channel, format!("msgid={before_msgid}"), 50) + .await + .context("Failed to send history before")?; + } } Ok(()) diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index a9597d6..6ab58f3 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -659,6 +659,19 @@ impl Tags { out } + #[cfg(feature = "web")] + pub fn server_time_with_fallback(&self) -> i64 { + self.server_time + .map(|t| t.unix_timestamp()) + .unwrap_or_else(|| { + web_time::SystemTime::now() + .duration_since(web_time::UNIX_EPOCH) + .unwrap() + .as_secs() as i64 + }) + } + + #[cfg(not(feature = "web"))] pub fn server_time_with_fallback(&self) -> i64 { self.server_time .unwrap_or_else(OffsetDateTime::now_utc) diff --git a/packages/core/core-wasm/Cargo.toml b/packages/core/core-wasm/Cargo.toml index 5ed5938..2459cb3 100644 --- a/packages/core/core-wasm/Cargo.toml +++ b/packages/core/core-wasm/Cargo.toml @@ -22,7 +22,7 @@ gloo-console = { version = "0.4.0" } gloo-net = "0.7.0" irc-proto = "1.1.0" -core-shared = { version = "0.1.0", path = "../core-shared", features = ["web"] } +core-shared = { version = "0.1.0", path = "../core-shared", features = ["web"], default-features = false } indexed_db_futures = { version = "0.6.4", features = ["async-upgrade", "serde"] } serde = { version = "1.0.228", features = ["derive"] } anyhow = { version = "1.0.103", features = ["backtrace"] } From 078d804c114a893b7f7baeb2201504301b461620 Mon Sep 17 00:00:00 2001 From: Jokler Date: Wed, 22 Jul 2026 01:50:15 +0200 Subject: [PATCH 05/28] Remove target requirment on CommandKey::History for now If no replies are sent there is no way to know which target a batch belongs to except for assuming that it's the next batch. Requested messages were reduced to 5 to simplify testing. --- packages/core/core-shared/src/actor.rs | 19 +++++++------------ packages/core/core-wasm/src/lib.rs | 14 +++++++------- 2 files changed, 14 insertions(+), 19 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 160d09c..6d73460 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -40,7 +40,7 @@ pub enum CommandKey { SignIn, Join(String), Privmsg { target: String, text: String }, - History { target: String }, + History, } #[derive(Debug)] @@ -54,6 +54,7 @@ pub enum CommandResponse { } impl ResponseChannels { + #[tracing::instrument] pub fn register(&mut self, key: CommandKey, os_tx: oneshot::Sender) { self.0.push((key, Instant::now(), os_tx)); } @@ -305,7 +306,7 @@ impl IrcActor { if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) && self.state.capabilities.history.enabled { - self.history_latest(channel_name.clone(), None, 50) + self.history_latest(channel_name.clone(), None, 5) .await .context("Failed to request latest history")?; } @@ -522,9 +523,7 @@ impl IrcActor { self.response_channels .reply( - &CommandKey::History { - target: batch.channel, - }, + &CommandKey::History, CommandResponse::History(history.clone()), ) .unwrap(); @@ -895,14 +894,10 @@ impl IrcActor { channel, before_msgid, } => { - self.response_channels.register( - CommandKey::History { - target: channel.clone(), - }, - cmd.reply_tx.unwrap(), - ); + self.response_channels + .register(CommandKey::History, cmd.reply_tx.unwrap()); - self.history_before(channel, format!("msgid={before_msgid}"), 50) + self.history_before(channel, format!("msgid={before_msgid}"), 5) .await .context("Failed to send history before")?; } diff --git a/packages/core/core-wasm/src/lib.rs b/packages/core/core-wasm/src/lib.rs index 276d521..c770e51 100644 --- a/packages/core/core-wasm/src/lib.rs +++ b/packages/core/core-wasm/src/lib.rs @@ -137,7 +137,7 @@ impl IrcConnection { .await .context("Failed to send ActorMessage")?; - let resp = rx.await.context("Failed to await ActorMessage")?; + let resp = rx.await.context("Failed to await actor state message")?; let CommandResponse::GetState(server) = resp else { unreachable!("expected state, got: {:?}", resp); }; @@ -253,7 +253,7 @@ impl IrcConnection { .await .context("Failed to send ActorMessage")?; - let resp = rx.await.context("Failed to await ActorMessage")?; + let resp = rx.await.context("Failed to await actor sign in message")?; let CommandResponse::SignIn(result) = resp else { unreachable!("expected sign in, got: {:?}", resp); }; @@ -281,7 +281,7 @@ impl IrcConnection { .await .context("Failed to send ActorMessage")?; - let resp = rx.await.context("Failed to await ActorMessage")?; + let resp = rx.await.context("Failed to await actor sign in message")?; let CommandResponse::SignIn(result) = resp else { unreachable!("expected sign in, got: {:?}", resp); @@ -305,7 +305,7 @@ impl IrcConnection { .await .context("Failed to send ActorMessage")?; - let resp = rx.await.context("Failed to await ActorMessage")?; + let resp = rx.await.context("Failed to await actor join message")?; let CommandResponse::Join(name) = resp else { unreachable!("expected join, got: {:?}", resp); }; @@ -334,7 +334,7 @@ impl IrcConnection { .await .context("Failed to send ActorMessage")?; - let resp = rx.await.context("Failed to await ActorMessage")?; + let resp = rx.await.context("Failed to await actor history message")?; let CommandResponse::History(history) = resp else { unreachable!("expected history, got: {:?}", resp); }; @@ -365,9 +365,9 @@ impl IrcChannel { .await .context("Failed to send ActorMessage")?; - let resp = rx.await.context("Failed to await ActorMessage")?; + let resp = rx.await.context("Failed to await actor message message")?; let CommandResponse::Privmsg(message) = resp else { - unreachable!("expected join, got: {:?}", resp); + unreachable!("expected privmsg, got: {:?}", resp); }; Ok((*message).into()) From 5b7a3024d3b83ef78f47da8c4d3f3815fa539a07 Mon Sep 17 00:00:00 2001 From: Jokler Date: Wed, 22 Jul 2026 01:51:27 +0200 Subject: [PATCH 06/28] Create a demo that connects to one channel and can request more history --- .../app/src/components/windows/WindowChat.vue | 52 ++++++++++++++- packages/app/src/stores/irc.ts | 63 ++++++++++++++++--- 2 files changed, 106 insertions(+), 9 deletions(-) diff --git a/packages/app/src/components/windows/WindowChat.vue b/packages/app/src/components/windows/WindowChat.vue index fedd234..dbe7f24 100644 --- a/packages/app/src/components/windows/WindowChat.vue +++ b/packages/app/src/components/windows/WindowChat.vue @@ -1,14 +1,28 @@ + diff --git a/packages/app/src/stores/irc.ts b/packages/app/src/stores/irc.ts index 592ec0d..b907f4b 100644 --- a/packages/app/src/stores/irc.ts +++ b/packages/app/src/stores/irc.ts @@ -1,6 +1,6 @@ import { defineStore } from "pinia" -import { Message, React, type IrcConnection, type Server, type ServerList } from "core-wasm" -import { computed, ref, shallowRef } from "vue" +import { ChannelMessage, Message, React, type IrcConnection, type Server, type ServerList, History, OrbitError, IrcChannel } from "core-wasm" +import { computed, reactive, ref, shallowRef } from "vue" import { useUserStore } from "./user" import { useAppStateStore } from "./app-state" @@ -22,7 +22,8 @@ export const useIrcStore = defineStore("irc", () => { const serverHandlers = shallowRef>(new Map()) // Holds references to messages per server. This should be actually per `server:channel` - const serverMessages = shallowRef>(new Map()) + const serverMessages = reactive>>(new Map()) + const serverChannel = ref() let controller: ServerList = {} as ServerList @@ -69,6 +70,7 @@ export const useIrcStore = defineStore("irc", () => { serverHandlers.value.set(state.id, handler) await handler.sign_in_anonymous(user.me.displayName, user.me.accountName, user.me.accountName) + serverChannel.value = await handler.join_channel("#orbit/testing") console.log("Signed in") registerServerEvents(state.id, handler) @@ -82,13 +84,29 @@ export const useIrcStore = defineStore("irc", () => { function registerServerEvents(key: number, handler: IrcConnection) { // Runs whenever some dataset on the server object changes handler.on_data((event) => { - if (event instanceof Message) { - const existing = serverMessages.value.get(key) ?? [] - existing.push(event) - serverMessages.value.set(key, existing) + if (event instanceof ChannelMessage) { + const existingServer = serverMessages.get(key) ?? new Map() + const existingChannel = existingServer.get(event.channel) ?? [] + + existingChannel.push(event.message) + existingChannel.sort((a: Message, b: Message) => a.metadata.server_time - b.metadata.server_time) + + existingServer.set(event.channel, existingChannel) + serverMessages.set(key, existingServer) } else if (event instanceof React) { // TODO console.log("Received reaction", event) + } else if (event instanceof History) { + const existingServer = serverMessages.get(key) ?? new Map() + const existingChannel = existingServer.get(event.channel) ?? [] + + for (const message of event.messages) { + existingChannel.push(message) + } + existingChannel.sort((a: Message, b: Message) => a.metadata.server_time - b.metadata.server_time) + + existingServer.set(event.channel, existingChannel) + serverMessages.set(key, existingServer) } }) @@ -108,6 +126,34 @@ export const useIrcStore = defineStore("irc", () => { return computed(() => serverState.value.get(id)) } + function getChannelMessages(id: number, channel: string) { + return computed(() => serverMessages.get(id)?.get(channel)) + } + + async function requestScrollback(id: number, channel: string) { + try { + const history = await serverHandlers.value.get(id)?.history_before(channel, serverMessages.get(id)?.get(channel)?.at(0)?.metadata.msgid ?? "") + + if (!!!history) { + return + } + + const existingServer = serverMessages.get(id) ?? new Map() + const existingChannel = existingServer.get(history.channel) ?? [] + + for (const message of history.messages) { + existingChannel.push(message) + } + existingChannel.sort((a: Message, b: Message) => a.metadata.server_time - b.metadata.server_time) + + existingServer.set(history.channel, existingChannel) + serverMessages.set(id, existingServer) + } catch (e: unknown) { + const error = e as OrbitError + console.error(JSON.parse(error.toString())) + } + } + return { init, serverConnect, @@ -116,5 +162,8 @@ export const useIrcStore = defineStore("irc", () => { serverData: serverState, serverControllers: serverHandlers, getServerState, + getChannelMessages, + requestScrollback, + serverChannel, } }) From e33a9e0fa03c373a87949fe09755a10cba05ecf7 Mon Sep 17 00:00:00 2001 From: Jokler Date: Thu, 23 Jul 2026 02:55:22 +0200 Subject: [PATCH 07/28] Associate history batch by label Other response channels should also use labels from now on if possible --- packages/core/Cargo.lock | 16 +++ packages/core/core-shared/Cargo.toml | 1 + packages/core/core-shared/src/actor.rs | 104 ++++++++++++++---- packages/core/core-shared/src/send_command.rs | 38 ++++--- packages/core/core-shared/src/state.rs | 20 ++-- 5 files changed, 136 insertions(+), 43 deletions(-) diff --git a/packages/core/Cargo.lock b/packages/core/Cargo.lock index 14daebf..10902ec 100644 --- a/packages/core/Cargo.lock +++ b/packages/core/Cargo.lock @@ -131,6 +131,7 @@ dependencies = [ "gloo-timers 0.4.0", "irc-proto", "ordermap", + "rand", "thiserror 2.0.18", "time", "tokio", @@ -851,6 +852,21 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "rand" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +dependencies = [ + "rand_core", +] + +[[package]] +name = "rand_core" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" + [[package]] name = "rustc_version" version = "0.4.1" diff --git a/packages/core/core-shared/Cargo.toml b/packages/core/core-shared/Cargo.toml index b34f2ba..bebd5ed 100644 --- a/packages/core/core-shared/Cargo.toml +++ b/packages/core/core-shared/Cargo.toml @@ -22,6 +22,7 @@ gloo-console = { version = "0.4.0", optional = true } gloo-net = { version = "0.7.0", optional = true } tokio = { version = "1.52.3", features = ["macros", "rt", "time", "sync"], optional = true } gloo-timers = { version = "0.4.0", features = ["futures"], optional = true } +rand = { version = "0.10.2", default-features = false } [features] web = ["dep:gloo-net", "dep:wasm-bindgen", "dep:gloo-console", "dep:tsify", "dep:web-time", "dep:gloo-timers"] diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 6d73460..e744070 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -3,6 +3,7 @@ use std::pin::pin; use std::{fmt, time::Duration}; use futures::FutureExt; +use rand::{SeedableRng, rngs::SmallRng, seq::IndexedRandom}; #[cfg(not(feature = "web"))] use std::time::Instant; #[cfg(feature = "web")] @@ -31,8 +32,11 @@ use irc_proto::{BatchSubCommand, CapSubCommand, Command::*, Message as IrcMessag use ordermap::OrderMap; use tracing::{debug, error, warn}; -#[derive(Debug, Default)] -pub struct ResponseChannels(Vec<(CommandKey, Instant, oneshot::Sender)>); +#[derive(Debug)] +pub struct ResponseChannels { + channels: Vec<(CommandKey, Instant, oneshot::Sender)>, + rng: SmallRng, +} #[derive(Debug, Clone, PartialEq, Eq)] pub enum CommandKey { @@ -41,6 +45,7 @@ pub enum CommandKey { Join(String), Privmsg { target: String, text: String }, History, + Label(String), } #[derive(Debug)] @@ -53,10 +58,45 @@ pub enum CommandResponse { History(History), } +const LABEL_CHARSET: &str = "abcdefghijklmnopqrstuvwxyz\ + ABCDEFGHIJKLMNOPQRSTUVWXYZ\ + 1234567890"; + impl ResponseChannels { + #[tracing::instrument] + pub fn new() -> Self { + Self { + channels: Vec::new(), + rng: SmallRng::from_seed([0; 32]), + } + } + #[tracing::instrument] pub fn register(&mut self, key: CommandKey, os_tx: oneshot::Sender) { - self.0.push((key, Instant::now(), os_tx)); + self.channels.push((key, Instant::now(), os_tx)); + } + + #[tracing::instrument] + pub fn register_labeled(&mut self, os_tx: oneshot::Sender) -> String { + let char_vec = LABEL_CHARSET + .split("") + .filter(|c| !c.is_empty()) + .collect::>(); + + let label = std::iter::repeat_with(|| { + char_vec + .choose(&mut self.rng) + .expect("CHARSET is not empty") + }) + .take(10) + .copied() + .collect::>() + .join(""); + + self.channels + .push((CommandKey::Label(label.clone()), Instant::now(), os_tx)); + + label } #[tracing::instrument] @@ -64,19 +104,24 @@ impl ResponseChannels { &mut self, key: &CommandKey, response: CommandResponse, - ) -> Result<(), CommandResponse> { - if let Some(idx) = self.0.iter().position(|(rk, _, _)| rk == key) { - let (_, _, ch) = self.0.remove(idx); + ) -> Result { + if let CommandKey::Label(label) = key { + dbg!(label); + } + + if let Some(idx) = self.channels.iter().position(|(rk, _, _)| rk == key) { + let (_, _, ch) = self.channels.remove(idx); ch.send(response)?; + + Ok(true) } else { // warn!("Failed to find response channel"); + Ok(false) } - - Ok(()) } pub fn check_timeouts(&mut self) { - self.0 + self.channels .retain(|(_, creation, _)| creation.elapsed() < Duration::from_secs(1)); } } @@ -136,6 +181,7 @@ pub trait IrcConnection: fmt::Debug { struct Batch { id: String, typ: BatchSubCommand, + label: Option, channel: String, messages: Vec, } @@ -188,7 +234,7 @@ impl IrcActor { incoming, outgoing, state: Server::new(id, address), - response_channels: ResponseChannels::default(), + response_channels: ResponseChannels::new(), event_handlers: Vec::new(), error_handlers: Vec::new(), disconnect_handlers: Vec::new(), @@ -306,7 +352,7 @@ impl IrcActor { if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) && self.state.capabilities.history.enabled { - self.history_latest(channel_name.clone(), None, 5) + self.history_latest(channel_name.clone(), None, 5, None) .await .context("Failed to request latest history")?; } @@ -454,7 +500,7 @@ impl IrcActor { return Ok(()); } - if let Some(username) = tags.username.clone() { + if let Some(username) = tags.account.clone() { let user = self.user_mut(source.to_string()).await; user.username = Some(username); } @@ -482,9 +528,15 @@ impl IrcActor { } BATCH(reference, typ, param) => { if let Some(id) = reference.strip_prefix('+') { + let mut tags = Tags::default(); + if let Some(ref t) = message.tags { + tags = Tags::parse(t); + } + self.current_batch = Some(Batch { id: id.to_string(), typ: typ.clone().unwrap(), + label: tags.label, channel: String::new(), messages: Vec::new(), }); @@ -521,13 +573,15 @@ impl IrcActor { messages: batch.messages, }; + let key = if let Some(label) = batch.label { + CommandKey::Label(label) + } else { + CommandKey::History + }; + self.response_channels - .reply( - &CommandKey::History, - CommandResponse::History(history.clone()), - ) + .reply(&key, CommandResponse::History(history.clone())) .unwrap(); - self.on_event(ServerEvent::History(history)).await?; } @@ -894,10 +948,19 @@ impl IrcActor { channel, before_msgid, } => { - self.response_channels - .register(CommandKey::History, cmd.reply_tx.unwrap()); + let label = if self.state.capabilities.labeled_response.enabled { + Some( + self.response_channels + .register_labeled(cmd.reply_tx.unwrap()), + ) + } else { + self.response_channels + .register(CommandKey::History, cmd.reply_tx.unwrap()); + + None + }; - self.history_before(channel, format!("msgid={before_msgid}"), 5) + self.history_before(channel, format!("msgid={before_msgid}"), 5, label) .await .context("Failed to send history before")?; } @@ -941,6 +1004,7 @@ impl IrcActor { .context("Failed to send CAPS LS")?; self.cap_req(&[ "echo-message", + "labeled-response", "message-tags", "sasl", "draft/message-redaction", diff --git a/packages/core/core-shared/src/send_command.rs b/packages/core/core-shared/src/send_command.rs index 9545f37..57194b7 100644 --- a/packages/core/core-shared/src/send_command.rs +++ b/packages/core/core-shared/src/send_command.rs @@ -4,7 +4,7 @@ use futures::{ SinkExt, channel::mpsc::{self, UnboundedSender}, }; -use irc_proto::{CapSubCommand, Command::*, Message as IrcMessage}; +use irc_proto::{CapSubCommand, Command::*, Message as IrcMessage, message::Tag}; pub trait SendCommand { type Error: std::error::Error + Send + Sync + 'static; @@ -14,10 +14,11 @@ pub trait SendCommand { fn command( &mut self, command: irc_proto::Command, + label: Option, ) -> impl Future> { async { self.message(IrcMessage { - tags: None, + tags: label.map(|l| vec![Tag(String::from("label"), Some(l))]), prefix: None, command, }) @@ -32,7 +33,7 @@ pub trait SendCommand { server1: String, server2: Option, ) -> impl std::future::Future> { - async { self.command(PONG(server1, server2)).await } + async { self.command(PONG(server1, server2), None).await } } fn cap_ls( @@ -40,7 +41,7 @@ pub trait SendCommand { version: String, ) -> impl std::future::Future> { async { - self.command(CAP(None, CapSubCommand::LS, Some(version), None)) + self.command(CAP(None, CapSubCommand::LS, Some(version), None), None) .await } } @@ -50,24 +51,27 @@ pub trait SendCommand { caps: &[&str], ) -> impl std::future::Future> { async { - self.command(CAP(None, CapSubCommand::REQ, None, Some(caps.join(" ")))) - .await + self.command( + CAP(None, CapSubCommand::REQ, None, Some(caps.join(" "))), + None, + ) + .await } } fn cap_end(&mut self) -> impl std::future::Future> { async { - self.command(CAP(None, CapSubCommand::END, None, None)) + self.command(CAP(None, CapSubCommand::END, None, None), None) .await } } fn nick(&mut self, nick: String) -> impl std::future::Future> { - async { self.command(NICK(nick)).await } + async { self.command(NICK(nick), None).await } } fn sasl(&mut self, req: String) -> impl std::future::Future> { - async { self.command(AUTHENTICATE(req)).await } + async { self.command(AUTHENTICATE(req), None).await } } fn sasl_plain(&mut self) -> impl std::future::Future> { @@ -84,7 +88,7 @@ pub trait SendCommand { mode: String, realname: String, ) -> impl std::future::Future> { - async { self.command(USER(user, mode, realname)).await } + async { self.command(USER(user, mode, realname), None).await } } fn join( @@ -92,7 +96,7 @@ pub trait SendCommand { channel: String, password: Option, ) -> impl std::future::Future> { - async { self.command(JOIN(channel, password, None)).await } + async { self.command(JOIN(channel, password, None), None).await } } fn privmsg( @@ -100,7 +104,7 @@ pub trait SendCommand { target: String, message: String, ) -> impl std::future::Future> { - async { self.command(PRIVMSG(target, message)).await } + async { self.command(PRIVMSG(target, message), None).await } } fn whois( @@ -108,17 +112,19 @@ pub trait SendCommand { server: Option, user: String, ) -> impl std::future::Future> { - async { self.command(WHOIS(server, user)).await } + async { self.command(WHOIS(server, user), None).await } } fn history( &mut self, subcommand: ChatHistorySubCommand, mut args: Vec, + label: Option, ) -> impl std::future::Future> { async move { args.insert(0, subcommand.to_string()); - self.command(Raw("CHATHISTORY".to_string(), args)).await + self.command(Raw("CHATHISTORY".to_string(), args), label) + .await } } @@ -127,11 +133,13 @@ pub trait SendCommand { target: String, before: String, limit: i32, + label: Option, ) -> impl std::future::Future> { async move { self.history( ChatHistorySubCommand::Before, vec![target, before, limit.to_string()], + label, ) .await } @@ -142,6 +150,7 @@ pub trait SendCommand { target: String, since: Option, limit: i32, + label: Option, ) -> impl std::future::Future> { async move { self.history( @@ -151,6 +160,7 @@ pub trait SendCommand { since.unwrap_or_else(|| String::from("*")), limit.to_string(), ], + label, ) .await } diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index 6ab58f3..fc76f1c 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -618,16 +618,17 @@ impl From for OrbitError { #[derive(Debug, Default)] pub struct Tags { - pub msgid: Option, pub server_time: Option, - pub username: Option, + pub msgid: Option, + pub account: Option, pub relayed_by: Option, + pub batch: Option, + pub bot: Option, + pub label: Option, pub reply: Option, pub react: Option, pub unreact: Option, - pub batch: Option, pub typing: Option, - pub bot: Option, } impl Tags { @@ -636,16 +637,17 @@ impl Tags { for Tag(key, value) in tags { match key.as_str() { - "msgid" => out.msgid = value.clone(), - "account" => out.username = value.clone(), - "draft/relaymsg" => out.relayed_by = value.clone(), - "batch" => out.batch = value.clone(), - "bot" => out.bot = value.clone(), "time" => { out.server_time = value .as_ref() .and_then(|v| OffsetDateTime::parse(v, &Iso8601::DEFAULT).ok()) } + "msgid" => out.msgid = value.clone(), + "account" => out.account = value.clone(), + "draft/relaymsg" => out.relayed_by = value.clone(), + "batch" => out.batch = value.clone(), + "bot" => out.bot = value.clone(), + "label" => out.label = dbg!(value.clone()), "+draft/reply" | "+reply" => out.reply = value.clone(), "+draft/react" => out.react = value.clone(), "+draft/unreact" => out.unreact = value.clone(), From 545c5ff8be23996db23b17932957e6b581faa99e Mon Sep 17 00:00:00 2001 From: Jokler Date: Thu, 23 Jul 2026 03:06:06 +0200 Subject: [PATCH 08/28] Rename ChannelRole::None to Regular to match standard IRC terms --- packages/core/core-shared/src/actor.rs | 4 ++-- packages/core/core-shared/src/state.rs | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index e744070..b29d3a6 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -362,7 +362,7 @@ impl IrcActor { channel.users.push(ChannelUser { nickname: source.to_string(), - role: ChannelRole::None, + role: ChannelRole::Regular, }); } } @@ -844,7 +844,7 @@ impl IrcActor { }); } else { channel_users.push(ChannelUser { - role: ChannelRole::None, + role: ChannelRole::Regular, nickname: user.clone(), }); } diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index fc76f1c..c429aad 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -469,7 +469,7 @@ pub enum ChannelRole { Operator, HalfOperator, Voice, - None, + Regular, } impl From for ChannelRole { From bc8bf48fcf5419ad8df77eaa9b0edef1951a9041 Mon Sep 17 00:00:00 2001 From: Jokler Date: Fri, 24 Jul 2026 03:18:28 +0200 Subject: [PATCH 09/28] Always track channel for requested chat history --- packages/core/core-shared/src/actor.rs | 117 ++++++++++++++++--------- 1 file changed, 74 insertions(+), 43 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index b29d3a6..1d59add 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -32,12 +32,6 @@ use irc_proto::{BatchSubCommand, CapSubCommand, Command::*, Message as IrcMessag use ordermap::OrderMap; use tracing::{debug, error, warn}; -#[derive(Debug)] -pub struct ResponseChannels { - channels: Vec<(CommandKey, Instant, oneshot::Sender)>, - rng: SmallRng, -} - #[derive(Debug, Clone, PartialEq, Eq)] pub enum CommandKey { RequestCaps, @@ -62,15 +56,36 @@ const LABEL_CHARSET: &str = "abcdefghijklmnopqrstuvwxyz\ ABCDEFGHIJKLMNOPQRSTUVWXYZ\ 1234567890"; -impl ResponseChannels { +fn generate_label(rng: &mut SmallRng) -> String { + let char_vec = LABEL_CHARSET + .split("") + .filter(|c| !c.is_empty()) + .collect::>(); + + std::iter::repeat_with(|| char_vec.choose(rng).expect("CHARSET is not empty")) + .take(10) + .copied() + .collect::>() + .join("") +} + +#[derive(Debug)] +pub struct ResponseChannels { + channels: Vec<(CommandKey, Instant, oneshot::Sender)>, + rng: SmallRng, +} + +impl Default for ResponseChannels { #[tracing::instrument] - pub fn new() -> Self { + fn default() -> Self { Self { channels: Vec::new(), rng: SmallRng::from_seed([0; 32]), } } +} +impl ResponseChannels { #[tracing::instrument] pub fn register(&mut self, key: CommandKey, os_tx: oneshot::Sender) { self.channels.push((key, Instant::now(), os_tx)); @@ -78,20 +93,7 @@ impl ResponseChannels { #[tracing::instrument] pub fn register_labeled(&mut self, os_tx: oneshot::Sender) -> String { - let char_vec = LABEL_CHARSET - .split("") - .filter(|c| !c.is_empty()) - .collect::>(); - - let label = std::iter::repeat_with(|| { - char_vec - .choose(&mut self.rng) - .expect("CHARSET is not empty") - }) - .take(10) - .copied() - .collect::>() - .join(""); + let label = generate_label(&mut self.rng); self.channels .push((CommandKey::Label(label.clone()), Instant::now(), os_tx)); @@ -177,8 +179,13 @@ pub trait IrcConnection: fmt::Debug { fn address(&self) -> &str; } +struct RequestedHistory { + channel: String, + label: Option, +} + #[derive(Debug)] -struct Batch { +struct CurrentBatch { id: String, typ: BatchSubCommand, label: Option, @@ -186,9 +193,9 @@ struct Batch { messages: Vec, } -impl Batch { +impl CurrentBatch { fn is_chathistory(&self) -> bool { - matches!(self.typ, BatchSubCommand::CUSTOM(ref c) if c.as_str() == "CHATHISTORY") + matches!(self.typ, BatchSubCommand::CUSTOM(ref c) if c.as_str() == "CHATHISTORY") } } @@ -214,8 +221,10 @@ pub struct IrcActor { error_handlers: Vec>, disconnect_handlers: Vec>, - current_batch: Option, + current_batch: Option, + requested_batches: Vec, sasl_state: SaslState, + rng: SmallRng, } impl IrcActor { @@ -234,12 +243,14 @@ impl IrcActor { incoming, outgoing, state: Server::new(id, address), - response_channels: ResponseChannels::new(), - event_handlers: Vec::new(), - error_handlers: Vec::new(), - disconnect_handlers: Vec::new(), - current_batch: None, + response_channels: ResponseChannels::default(), + event_handlers: Default::default(), + error_handlers: Default::default(), + disconnect_handlers: Default::default(), + current_batch: Default::default(), + requested_batches: Default::default(), sasl_state: Default::default(), + rng: SmallRng::from_seed([1; 32]), }; let (tx, rx) = oneshot::channel(); @@ -352,7 +363,17 @@ impl IrcActor { if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) && self.state.capabilities.history.enabled { - self.history_latest(channel_name.clone(), None, 5, None) + let label = if self.state.capabilities.labeled_response.enabled { + Some(generate_label(&mut self.rng)) + } else { + None + }; + + self.requested_batches.push(RequestedHistory { + channel: channel_name.clone(), + label: label.clone(), + }); + self.history_latest(channel_name.clone(), None, 5, label) .await .context("Failed to request latest history")?; } @@ -533,11 +554,21 @@ impl IrcActor { tags = Tags::parse(t); } - self.current_batch = Some(Batch { + let Some(idx) = self + .requested_batches + .iter() + .position(|b| b.label == tags.label) + else { + // TODO: handlee unrequested batches + return Ok(()); + }; + let channel = self.requested_batches.remove(idx).channel; + + self.current_batch = Some(CurrentBatch { id: id.to_string(), typ: typ.clone().unwrap(), label: tags.label, - channel: String::new(), + channel, messages: Vec::new(), }); @@ -547,19 +578,15 @@ impl IrcActor { _ => warn!(?typ, ?param, "unhandled BATCH type"), } } else { + if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { + // TODO: handle other types of batches + return Ok(()); + } assert_eq!( - self.current_batch.as_ref().map(|s| s.id.as_str()), + self.current_batch.as_ref().map(|b| b.id.as_str()), Some(&reference[1..]) ); - // FIXME: this will panic if a batch only contains QUITs - assert_eq!( - self.current_batch.as_ref().map(|b| b.is_chathistory() - && !b.messages.is_empty() - && b.messages.is_empty()), - Some(false) - ); - if let Some(batch) = self.current_batch.take() { for message in &batch.messages { let channel = self.channel_mut(batch.channel.clone()).await; @@ -960,6 +987,10 @@ impl IrcActor { None }; + self.requested_batches.push(RequestedHistory { + channel: channel.clone(), + label: label.clone(), + }); self.history_before(channel, format!("msgid={before_msgid}"), 5, label) .await .context("Failed to send history before")?; From 53070016fe3bba2e28b715a5442508a7c96e9615 Mon Sep 17 00:00:00 2001 From: dolanske Date: Mon, 27 Jul 2026 10:19:47 +0300 Subject: [PATCH 10/28] Remove duplicate history saving, clean up --- .../app/src/components/windows/WindowChat.vue | 4 --- packages/app/src/stores/irc.ts | 25 ++++++------------- 2 files changed, 7 insertions(+), 22 deletions(-) diff --git a/packages/app/src/components/windows/WindowChat.vue b/packages/app/src/components/windows/WindowChat.vue index dbe7f24..53e79c9 100644 --- a/packages/app/src/components/windows/WindowChat.vue +++ b/packages/app/src/components/windows/WindowChat.vue @@ -4,16 +4,12 @@ import type { WindowChat } from "../../lib/windows" import { useIrcStore } from "../../stores/irc" import { ref } from "vue" -// TODO: if no chat is active, list available chats for a server. We know that -// because channelId will be "__unspecified" - const props = defineProps() const editmsg = ref("") const irc = useIrcStore() const messages = irc.getChannelMessages(props.serverId, "#orbit/testing") -// const channel = irc.get async function sendMessage() { if (editmsg.value === "") { diff --git a/packages/app/src/stores/irc.ts b/packages/app/src/stores/irc.ts index b907f4b..5ad439f 100644 --- a/packages/app/src/stores/irc.ts +++ b/packages/app/src/stores/irc.ts @@ -1,5 +1,5 @@ import { defineStore } from "pinia" -import { ChannelMessage, Message, React, type IrcConnection, type Server, type ServerList, History, OrbitError, IrcChannel } from "core-wasm" +import { ChannelMessage, Message, React, type IrcConnection, type Server, type ServerList, OrbitError, IrcChannel } from "core-wasm" import { computed, reactive, ref, shallowRef } from "vue" import { useUserStore } from "./user" import { useAppStateStore } from "./app-state" @@ -96,17 +96,6 @@ export const useIrcStore = defineStore("irc", () => { } else if (event instanceof React) { // TODO console.log("Received reaction", event) - } else if (event instanceof History) { - const existingServer = serverMessages.get(key) ?? new Map() - const existingChannel = existingServer.get(event.channel) ?? [] - - for (const message of event.messages) { - existingChannel.push(message) - } - existingChannel.sort((a: Message, b: Message) => a.metadata.server_time - b.metadata.server_time) - - existingServer.set(event.channel, existingChannel) - serverMessages.set(key, existingServer) } }) @@ -122,6 +111,7 @@ export const useIrcStore = defineStore("irc", () => { }) } + // TODO: these should be cached not to create a separate computed value on each call function getServerState(id: number) { return computed(() => serverState.value.get(id)) } @@ -130,21 +120,20 @@ export const useIrcStore = defineStore("irc", () => { return computed(() => serverMessages.get(id)?.get(channel)) } + // TODO: will be called automatically by a scroll listener to append new messages as user's nearing the top of the window async function requestScrollback(id: number, channel: string) { try { const history = await serverHandlers.value.get(id)?.history_before(channel, serverMessages.get(id)?.get(channel)?.at(0)?.metadata.msgid ?? "") - if (!!!history) { + if (!history) { return } - const existingServer = serverMessages.get(id) ?? new Map() + const existingServer = serverMessages.get(id) ?? new Map() const existingChannel = existingServer.get(history.channel) ?? [] - for (const message of history.messages) { - existingChannel.push(message) - } - existingChannel.sort((a: Message, b: Message) => a.metadata.server_time - b.metadata.server_time) + existingChannel.push(...history.messages) + existingChannel.sort((a, b) => a.metadata.server_time - b.metadata.server_time) existingServer.set(history.channel, existingChannel) serverMessages.set(id, existingServer) From 08fae9c7a03fe12aeb1290eb7b7c868a42f5ad89 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 13:10:56 +0200 Subject: [PATCH 11/28] Track all batches and load initial channel history again This should fix the assert getting hit when an unrequested batch gets handeled. --- packages/app/src/stores/irc.ts | 12 +++ packages/core/core-shared/src/actor.rs | 140 +++++++++++++++---------- packages/core/core-wasm/src/lib.rs | 19 ++++ 3 files changed, 117 insertions(+), 54 deletions(-) diff --git a/packages/app/src/stores/irc.ts b/packages/app/src/stores/irc.ts index 5ad439f..46544e7 100644 --- a/packages/app/src/stores/irc.ts +++ b/packages/app/src/stores/irc.ts @@ -73,6 +73,18 @@ export const useIrcStore = defineStore("irc", () => { serverChannel.value = await handler.join_channel("#orbit/testing") console.log("Signed in") + // Set initial channel messages + const channelState = await serverChannel.value.state() + + const existingServer = serverMessages.get(0) ?? new Map() + const existingChannel = existingServer.get(channelState.metadata.name) ?? [] + + existingChannel.push(...channelState.messages) + existingChannel.sort((a, b) => a.metadata.server_time - b.metadata.server_time) + + existingServer.set(channelState.metadata.name, existingChannel) + serverMessages.set(state.id, existingServer) + registerServerEvents(state.id, handler) return { diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 1d59add..fa6ed16 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -45,6 +45,7 @@ pub enum CommandKey { #[derive(Debug)] pub enum CommandResponse { GetState(Box), + GetChannelState(Option), Capabilities, SignIn(Result), Join(String), @@ -137,6 +138,7 @@ pub struct ActorMessage { #[derive(Debug)] pub enum ActorCommand { GetState, + GetChannelState(String), SignIn { nick: String, user: String, @@ -187,15 +189,22 @@ struct RequestedHistory { #[derive(Debug)] struct CurrentBatch { id: String, - typ: BatchSubCommand, - label: Option, - channel: String, - messages: Vec, + data: BatchData, +} + +#[derive(Debug)] +enum BatchData { + History { + label: Option, + channel: String, + messages: Vec, + }, + Unhandled, } impl CurrentBatch { fn is_chathistory(&self) -> bool { - matches!(self.typ, BatchSubCommand::CUSTOM(ref c) if c.as_str() == "CHATHISTORY") + matches!(self.data, BatchData::History { .. }) } } @@ -222,7 +231,7 @@ pub struct IrcActor { disconnect_handlers: Vec>, current_batch: Option, - requested_batches: Vec, + requested_history_batches: Vec, sasl_state: SaslState, rng: SmallRng, } @@ -248,7 +257,7 @@ impl IrcActor { error_handlers: Default::default(), disconnect_handlers: Default::default(), current_batch: Default::default(), - requested_batches: Default::default(), + requested_history_batches: Default::default(), sasl_state: Default::default(), rng: SmallRng::from_seed([1; 32]), }; @@ -351,12 +360,15 @@ impl IrcActor { self.state .channels .insert(channel_name.clone(), channel.clone()); - self.response_channels - .reply( - &CommandKey::Join(channel_name.clone()), - CommandResponse::Join(channel_name.clone()), - ) - .map_err(|e| anyhow!("Failed to reply to JOIN command {e:?}"))?; + + if !self.state.capabilities.history.enabled { + self.response_channels + .reply( + &CommandKey::Join(channel_name.clone()), + CommandResponse::Join(channel_name.clone()), + ) + .map_err(|e| anyhow!("Failed to reply to JOIN command {e:?}"))?; + } self.on_event(ServerEvent::Joined(channel)).await?; @@ -369,7 +381,7 @@ impl IrcActor { None }; - self.requested_batches.push(RequestedHistory { + self.requested_history_batches.push(RequestedHistory { channel: channel_name.clone(), label: label.clone(), }); @@ -448,9 +460,9 @@ impl IrcActor { }; if let Some(batch) = self.current_batch.as_mut() - && batch.is_chathistory() + && let BatchData::History { messages, .. } = &mut batch.data { - batch.messages.push(state_message); + messages.push(state_message); return Ok(()); } @@ -554,62 +566,73 @@ impl IrcActor { tags = Tags::parse(t); } - let Some(idx) = self - .requested_batches - .iter() - .position(|b| b.label == tags.label) - else { - // TODO: handlee unrequested batches - return Ok(()); - }; - let channel = self.requested_batches.remove(idx).channel; - - self.current_batch = Some(CurrentBatch { - id: id.to_string(), - typ: typ.clone().unwrap(), - label: tags.label, - channel, - messages: Vec::new(), - }); - match typ { - Some(BatchSubCommand::CUSTOM(c)) - if &c == "METADATA" || &c == "CHATHISTORY" => {} - _ => warn!(?typ, ?param, "unhandled BATCH type"), + Some(BatchSubCommand::CUSTOM(c)) if &c == "CHATHISTORY" => { + let idx = self + .requested_history_batches + .iter() + .position(|b| b.label == tags.label) + .expect("Chat history was requested"); + let channel = self.requested_history_batches.remove(idx).channel; + + self.current_batch = Some(CurrentBatch { + id: id.to_string(), + data: BatchData::History { + label: tags.label, + channel, + messages: Vec::new(), + }, + }); + } + _ => { + self.current_batch = Some(CurrentBatch { + id: id.to_string(), + data: BatchData::Unhandled, + }); + warn!(?typ, ?param, "unhandled BATCH type"); + } } } else { - if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { - // TODO: handle other types of batches - return Ok(()); - } assert_eq!( self.current_batch.as_ref().map(|b| b.id.as_str()), Some(&reference[1..]) ); - if let Some(batch) = self.current_batch.take() { - for message in &batch.messages { - let channel = self.channel_mut(batch.channel.clone()).await; + if let Some(batch) = self.current_batch.take() + && let BatchData::History { + label, + channel: channel_name, + messages, + } = batch.data + { + let channel = self.channel_mut(channel_name.clone()).await; + for message in &messages { channel .messages .insert(message.metadata.msgid.clone(), message.clone()); } let history = History { - channel: batch.channel.clone(), - messages: batch.messages, + channel: channel_name.clone(), + messages, }; - let key = if let Some(label) = batch.label { + let key = if let Some(label) = label { CommandKey::Label(label) } else { CommandKey::History }; + self.response_channels + .reply( + &CommandKey::Join(channel_name.clone()), + CommandResponse::Join(channel_name.clone()), + ) + .map_err(|e| anyhow!("Failed to reply to JOIN command {e:?}"))?; + self.response_channels .reply(&key, CommandResponse::History(history.clone())) .unwrap(); - self.on_event(ServerEvent::History(history)).await?; } self.current_batch = None; @@ -727,14 +750,16 @@ impl IrcActor { pub async fn push_batch(&mut self, target: String, state_message: Message) -> bool { if let Some(batch) = self.current_batch.as_mut() - && batch.is_chathistory() + && let BatchData::History { + channel, messages, .. + } = &mut batch.data { - if batch.channel.is_empty() { - batch.channel = target.clone(); + if channel.is_empty() { + *channel = target.clone(); } else { - assert_eq!(batch.channel, target) + assert_eq!(*channel, target) } - batch.messages.push(state_message); + messages.push(state_message); return true; } @@ -928,6 +953,13 @@ impl IrcActor { .unwrap() .send(CommandResponse::GetState(Box::new(self.state.clone()))) .unwrap(), + ActorCommand::GetChannelState(channel_name) => cmd + .reply_tx + .unwrap() + .send(CommandResponse::GetChannelState( + self.state.channels.get(&channel_name).cloned(), + )) + .unwrap(), ActorCommand::SignIn { nick, user, @@ -987,7 +1019,7 @@ impl IrcActor { None }; - self.requested_batches.push(RequestedHistory { + self.requested_history_batches.push(RequestedHistory { channel: channel.clone(), label: label.clone(), }); diff --git a/packages/core/core-wasm/src/lib.rs b/packages/core/core-wasm/src/lib.rs index c770e51..49fa205 100644 --- a/packages/core/core-wasm/src/lib.rs +++ b/packages/core/core-wasm/src/lib.rs @@ -351,6 +351,25 @@ pub struct IrcChannel { #[wasm_bindgen] impl IrcChannel { + #[wasm_bindgen] + pub async fn state(&mut self) -> Result { + let (tx, rx) = oneshot::channel(); + self.address + .send(ActorMessage { + command: ActorCommand::GetChannelState(self.name.clone()), + reply_tx: Some(tx), + }) + .await + .context("Failed to send ActorMessage")?; + + let resp = rx.await.context("Failed to await actor state message")?; + let CommandResponse::GetChannelState(channel) = resp else { + unreachable!("expected state, got: {:?}", resp); + }; + + Ok(channel.unwrap().into()) + } + #[wasm_bindgen] pub async fn send_message(&mut self, text: String) -> Result { let (tx, rx) = oneshot::channel(); From ac761de75a1168dd20357c61d9035949bcd0e5e5 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 13:14:09 +0200 Subject: [PATCH 12/28] Fix clippy warning --- packages/core/core-shared/src/actor.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index fa6ed16..76fd820 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -45,7 +45,7 @@ pub enum CommandKey { #[derive(Debug)] pub enum CommandResponse { GetState(Box), - GetChannelState(Option), + GetChannelState(Box>), Capabilities, SignIn(Result), Join(String), @@ -956,9 +956,9 @@ impl IrcActor { ActorCommand::GetChannelState(channel_name) => cmd .reply_tx .unwrap() - .send(CommandResponse::GetChannelState( + .send(CommandResponse::GetChannelState(Box::new( self.state.channels.get(&channel_name).cloned(), - )) + ))) .unwrap(), ActorCommand::SignIn { nick, From 0f9ecf0277a48e0e3f75732bf51e26c27cf46a60 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 13:21:33 +0200 Subject: [PATCH 13/28] Set MessageReference fields to None if the Message was not found --- packages/core/core-shared/src/actor.rs | 9 ++------- packages/core/core-shared/src/state.rs | 4 +++- 2 files changed, 5 insertions(+), 8 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 76fd820..95591b3 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -503,13 +503,8 @@ impl IrcActor { .and_then(|c| c.messages.get(r)) }) .map(|m| MessageReference { - text: m - .map(|m| m.text.clone().map(|t| t.content)) - .unwrap_or_else(|| Some(String::from("Unknown"))), - - username: m - .map(|m| m.metadata.user.clone()) - .unwrap_or_else(|| String::from("Unknown")), + text: m.and_then(|m| m.text.clone().map(|t| t.content)), + username: m.map(|m| m.metadata.user.clone()), }); let state_message = Message { diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index c429aad..d6bdb71 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -550,8 +550,10 @@ impl Eq for MessageMetadata {} #[cfg_attr(feature = "web", derive(Tsify))] #[cfg_attr(feature = "web", wasm_bindgen(getter_with_clone, inspectable))] pub struct MessageReference { + /// Empty if message wasn't found or if reply wasn't to a text message pub text: Option, - pub username: String, + /// Empty if message wasn't found + pub username: Option, } #[derive(Debug, Clone, PartialEq, Eq)] From 37125cbb051565a4b62959e0a7fa46edc046aa4b Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 13:26:39 +0200 Subject: [PATCH 14/28] Add message type to msgid fallback and fix typos --- packages/core/core-shared/src/actor.rs | 10 +++++----- packages/core/core-wasm/src/lib.rs | 2 +- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 95591b3..6efa5e8 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -341,7 +341,7 @@ impl IrcActor { let state_message = Message { text: None, metadata: MessageMetadata { - msgid: tags.msgid_with_fallback(&[source]), + msgid: tags.msgid_with_fallback(&["JOIN", source]), server_time: tags.server_time_with_fallback() as f64, message_type: MessageType::Join, user: source.to_string(), @@ -416,7 +416,7 @@ impl IrcActor { let state_message = Message { text: None, metadata: MessageMetadata { - msgid: tags.msgid_with_fallback(&[source]), + msgid: tags.msgid_with_fallback(&["PART", source]), server_time: tags.server_time_with_fallback() as f64, message_type: MessageType::Part, user: source.to_string(), @@ -452,9 +452,9 @@ impl IrcActor { let state_message = Message { text: None, metadata: MessageMetadata { - msgid: tags.msgid_with_fallback(&[source]), + msgid: tags.msgid_with_fallback(&["QUIT", source]), server_time: tags.server_time_with_fallback() as f64, - message_type: MessageType::Part, + message_type: MessageType::Quit, user: source.to_string(), }, }; @@ -491,7 +491,7 @@ impl IrcActor { let source = message.source_nickname().unwrap(); - let msgid = tags.msgid_with_fallback(&[source, target, text]); + let msgid = tags.msgid_with_fallback(&["PRIVMSG", source, target, text]); let reply = tags .reply diff --git a/packages/core/core-wasm/src/lib.rs b/packages/core/core-wasm/src/lib.rs index 49fa205..f207ca3 100644 --- a/packages/core/core-wasm/src/lib.rs +++ b/packages/core/core-wasm/src/lib.rs @@ -384,7 +384,7 @@ impl IrcChannel { .await .context("Failed to send ActorMessage")?; - let resp = rx.await.context("Failed to await actor message message")?; + let resp = rx.await.context("Failed to await actor message")?; let CommandResponse::Privmsg(message) = resp else { unreachable!("expected privmsg, got: {:?}", resp); }; From 20873a097477fdaf513fddd344ae6b92fa02686c Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 13:32:28 +0200 Subject: [PATCH 15/28] Only recreate timeouts after they timeout --- packages/core/core-shared/src/actor.rs | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 6efa5e8..8e827f1 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -2,7 +2,7 @@ use std::pin::pin; use std::{fmt, time::Duration}; -use futures::FutureExt; +use futures::{FutureExt, future::FusedFuture}; use rand::{SeedableRng, rngs::SmallRng, seq::IndexedRandom}; #[cfg(not(feature = "web"))] use std::time::Instant; @@ -277,12 +277,18 @@ impl IrcActor { #[tracing::instrument(skip(self))] pub async fn run(mut self) { - loop { + fn create_timeout() -> impl FusedFuture { #[cfg(feature = "web")] let mut timeout = gloo_timers::future::TimeoutFuture::new(1000).fuse(); #[cfg(not(feature = "web"))] let mut timeout = pin!(tokio::time::sleep(Duration::from_secs(1)).fuse()); + timeout + } + + let mut timeout = create_timeout(); + + loop { futures::select! { msg = self.incoming.next() => { match msg { @@ -306,6 +312,8 @@ impl IrcActor { } _ = timeout => { self.response_channels.check_timeouts(); + + timeout = create_timeout(); } } } From 96016efd1376ce7beeb097ca4cd855109122bfb8 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 14:00:23 +0200 Subject: [PATCH 16/28] Remove superfluous chathistory checks --- packages/core/core-shared/src/actor.rs | 30 ++++++++++---------------- 1 file changed, 11 insertions(+), 19 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 8e827f1..f06a0e5 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -380,9 +380,7 @@ impl IrcActor { self.on_event(ServerEvent::Joined(channel)).await?; - if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) - && self.state.capabilities.history.enabled - { + if self.state.capabilities.history.enabled { let label = if self.state.capabilities.labeled_response.enabled { Some(generate_label(&mut self.rng)) } else { @@ -398,14 +396,12 @@ impl IrcActor { .context("Failed to request latest history")?; } } else { - if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { - let channel = self.channel_mut(channel_name.clone()).await; + let channel = self.channel_mut(channel_name.clone()).await; - channel.users.push(ChannelUser { - nickname: source.to_string(), - role: ChannelRole::Regular, - }); - } + channel.users.push(ChannelUser { + nickname: source.to_string(), + role: ChannelRole::Regular, + }); } self.on_event(ServerEvent::Privmsg { @@ -444,11 +440,9 @@ impl IrcActor { }) .await?; - if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { - let channel = self.channel_mut(channel_name.clone()).await; + let channel = self.channel_mut(channel_name.clone()).await; - channel.users.retain(|u| u.nickname != source); - } + channel.users.retain(|u| u.nickname != source); } QUIT(ref comment) => { let mut tags = Tags::default(); @@ -480,11 +474,9 @@ impl IrcActor { }) .await?; - if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { - self.state.users.remove(source); - for channel in self.state.channels.values_mut() { - channel.users.retain(|u| u.nickname != source); - } + self.state.users.remove(source); + for channel in self.state.channels.values_mut() { + channel.users.retain(|u| u.nickname != source); } } PRIVMSG(ref target, ref text) => { From 43ae29e42c224a21107feeb92bff9f758b07f87a Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 14:02:10 +0200 Subject: [PATCH 17/28] Fix reactor state update mix up --- packages/core/core-shared/src/actor.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index f06a0e5..7dbc5ad 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -674,9 +674,9 @@ impl IrcActor { .or_insert_with(Vec::new); if is_unreact { - reactors.push(nickname.clone()); - } else { reactors.retain(|v| *v != nickname); + } else { + reactors.push(nickname.clone()); } // TODO: should it be sent if the message wasn't found? From 653ce89f70bf6bd0080affa90725e5741ac65ca5 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 14:08:09 +0200 Subject: [PATCH 18/28] Ignore scrollback request if there are no messages --- packages/app/src/stores/irc.ts | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/packages/app/src/stores/irc.ts b/packages/app/src/stores/irc.ts index 46544e7..278fcbd 100644 --- a/packages/app/src/stores/irc.ts +++ b/packages/app/src/stores/irc.ts @@ -135,7 +135,11 @@ export const useIrcStore = defineStore("irc", () => { // TODO: will be called automatically by a scroll listener to append new messages as user's nearing the top of the window async function requestScrollback(id: number, channel: string) { try { - const history = await serverHandlers.value.get(id)?.history_before(channel, serverMessages.get(id)?.get(channel)?.at(0)?.metadata.msgid ?? "") + const oldestId = serverMessages.get(id)?.get(channel)?.at(0)?.metadata.msgid + if (!oldestId) { + return + } + const history = await serverHandlers.value.get(id)?.history_before(channel, oldestId) if (!history) { return From 13c77408133b2e8ab09829f6150a7fe3c525d59e Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 14:11:40 +0200 Subject: [PATCH 19/28] Turn debug statement into a warning log --- packages/core/core-shared/src/actor.rs | 7 +++---- packages/core/core-shared/src/state.rs | 2 +- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 7dbc5ad..33d0064 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -108,16 +108,15 @@ impl ResponseChannels { key: &CommandKey, response: CommandResponse, ) -> Result { - if let CommandKey::Label(label) = key { - dbg!(label); - } - if let Some(idx) = self.channels.iter().position(|(rk, _, _)| rk == key) { let (_, _, ch) = self.channels.remove(idx); ch.send(response)?; Ok(true) } else { + if let CommandKey::Label(label) = key { + warn!("Failed to find response channel for label {label:?}"); + } // warn!("Failed to find response channel"); Ok(false) } diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index d6bdb71..2cf809c 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -649,7 +649,7 @@ impl Tags { "draft/relaymsg" => out.relayed_by = value.clone(), "batch" => out.batch = value.clone(), "bot" => out.bot = value.clone(), - "label" => out.label = dbg!(value.clone()), + "label" => out.label = value.clone(), "+draft/reply" | "+reply" => out.reply = value.clone(), "+draft/react" => out.react = value.clone(), "+draft/unreact" => out.unreact = value.clone(), From 19fddc48847a7a9c87da84d9c628f7a2d7bbf080 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 14:48:02 +0200 Subject: [PATCH 20/28] Sync import cfg directive with use location --- packages/core/core-shared/src/actor.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 33d0064..ebe913e 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -1,4 +1,4 @@ -#[cfg(feature = "default")] +#[cfg(not(feature = "web"))] use std::pin::pin; use std::{fmt, time::Duration}; From 9a679d0de8ec46c714295d473e2cc1a5f7d6f832 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 4 Aug 2026 14:56:50 +0200 Subject: [PATCH 21/28] Remove unused History ServerEvent --- packages/core/core-shared/src/actor.rs | 4 ++-- packages/core/core-shared/src/state.rs | 1 - packages/core/core-wasm/src/lib.rs | 2 -- 3 files changed, 2 insertions(+), 5 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index ebe913e..ffb0b7a 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -278,9 +278,9 @@ impl IrcActor { pub async fn run(mut self) { fn create_timeout() -> impl FusedFuture { #[cfg(feature = "web")] - let mut timeout = gloo_timers::future::TimeoutFuture::new(1000).fuse(); + let timeout = gloo_timers::future::TimeoutFuture::new(1000).fuse(); #[cfg(not(feature = "web"))] - let mut timeout = pin!(tokio::time::sleep(Duration::from_secs(1)).fuse()); + let timeout = pin!(tokio::time::sleep(Duration::from_secs(1)).fuse()); timeout } diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index 2cf809c..abb0753 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -575,7 +575,6 @@ pub enum ServerEvent { text: String, is_unreact: bool, }, - History(History), } #[derive(Debug, Clone, PartialEq, Eq)] diff --git a/packages/core/core-wasm/src/lib.rs b/packages/core/core-wasm/src/lib.rs index f207ca3..981f988 100644 --- a/packages/core/core-wasm/src/lib.rs +++ b/packages/core/core-wasm/src/lib.rs @@ -520,7 +520,6 @@ pub enum ServerEvent { UserList(UserList), Privmsg(ChannelMessage), React(React), - History(History), } impl From for ServerEvent { @@ -547,7 +546,6 @@ impl From for ServerEvent { text, is_unreact, }), - state::ServerEvent::History(history) => Self::History(history.into()), } } } From dbad4e49a18d2ffab191a58b4068b8b173f282d5 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 11 Aug 2026 16:43:04 +0200 Subject: [PATCH 22/28] Always get initial server messages from correct server --- packages/app/src/stores/irc.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/app/src/stores/irc.ts b/packages/app/src/stores/irc.ts index 278fcbd..b6bd2a1 100644 --- a/packages/app/src/stores/irc.ts +++ b/packages/app/src/stores/irc.ts @@ -76,7 +76,7 @@ export const useIrcStore = defineStore("irc", () => { // Set initial channel messages const channelState = await serverChannel.value.state() - const existingServer = serverMessages.get(0) ?? new Map() + const existingServer = serverMessages.get(state.id) ?? new Map() const existingChannel = existingServer.get(channelState.metadata.name) ?? [] existingChannel.push(...channelState.messages) From a47b27a33dd80fd97d8268da257bb6b33fac9d52 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 11 Aug 2026 16:37:58 +0200 Subject: [PATCH 23/28] Fix no web compile error --- packages/core/core-shared/src/actor.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index ffb0b7a..381af51 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -1,5 +1,3 @@ -#[cfg(not(feature = "web"))] -use std::pin::pin; use std::{fmt, time::Duration}; use futures::{FutureExt, future::FusedFuture}; @@ -280,7 +278,7 @@ impl IrcActor { #[cfg(feature = "web")] let timeout = gloo_timers::future::TimeoutFuture::new(1000).fuse(); #[cfg(not(feature = "web"))] - let timeout = pin!(tokio::time::sleep(Duration::from_secs(1)).fuse()); + let timeout = Box::pin(tokio::time::sleep(Duration::from_secs(1)).fuse()); timeout } From 587857be5be0e3b8a155a0b4e9c5991e01e76a72 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 11 Aug 2026 18:50:54 +0200 Subject: [PATCH 24/28] Add support for nested batches and multiline messages --- packages/core/core-shared/src/actor.rs | 228 ++++++++++++++++++------- 1 file changed, 170 insertions(+), 58 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index 381af51..d6370ec 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -196,6 +196,10 @@ enum BatchData { channel: String, messages: Vec, }, + Multiline { + target: String, + message: Message, + }, Unhandled, } @@ -227,7 +231,7 @@ pub struct IrcActor { error_handlers: Vec>, disconnect_handlers: Vec>, - current_batch: Option, + current_batches: Vec, requested_history_batches: Vec, sasl_state: SaslState, rng: SmallRng, @@ -253,7 +257,7 @@ impl IrcActor { event_handlers: Default::default(), error_handlers: Default::default(), disconnect_handlers: Default::default(), - current_batch: Default::default(), + current_batches: Default::default(), requested_history_batches: Default::default(), sasl_state: Default::default(), rng: SmallRng::from_seed([1; 32]), @@ -339,7 +343,7 @@ impl IrcActor { tags = Tags::parse(t); } - if self.current_batch.as_ref().map(|b| b.is_chathistory()) == Some(true) { + if self.current_batches.iter().any(|b| b.is_chathistory()) { assert!(tags.server_time.is_some()); } @@ -458,11 +462,11 @@ impl IrcActor { }, }; - if let Some(batch) = self.current_batch.as_mut() - && let BatchData::History { messages, .. } = &mut batch.data - { - messages.push(state_message); - return Ok(()); + for batch in &mut self.current_batches { + if let BatchData::History { messages, .. } = &mut batch.data { + messages.push(state_message); + return Ok(()); + } } self.on_event(ServerEvent::Privmsg { @@ -482,10 +486,26 @@ impl IrcActor { tags = Tags::parse(t); } assert_eq!( - self.current_batch.as_ref().map(|s| s.id.as_str()), + self.current_batches.iter().last().map(|b| b.id.as_str()), tags.batch.as_deref(), ); + if let Some(batch) = self.current_batches.iter_mut().last() + && let BatchData::Multiline { + message, + target: channel, + } = &mut batch.data + && let Some(t) = message.text.as_mut() + { + if t.content.is_empty() { + *channel = target.to_string(); + t.content = text.to_string(); + } else { + t.content = format!("{}\n{}", t.content, text); + } + return Ok(()); + } + let source = message.source_nickname().unwrap(); let msgid = tags.msgid_with_fallback(&["PRIVMSG", source, target, text]); @@ -551,15 +571,14 @@ impl IrcActor { }) .await?; } - BATCH(reference, typ, param) => { + BATCH(ref reference, ref typ, ref param) => { + let mut tags = Tags::default(); + if let Some(ref t) = message.tags { + tags = Tags::parse(t); + } if let Some(id) = reference.strip_prefix('+') { - let mut tags = Tags::default(); - if let Some(ref t) = message.tags { - tags = Tags::parse(t); - } - match typ { - Some(BatchSubCommand::CUSTOM(c)) if &c == "CHATHISTORY" => { + Some(BatchSubCommand::CUSTOM(c)) if c.as_str() == "CHATHISTORY" => { let idx = self .requested_history_batches .iter() @@ -567,7 +586,7 @@ impl IrcActor { .expect("Chat history was requested"); let channel = self.requested_history_batches.remove(idx).channel; - self.current_batch = Some(CurrentBatch { + self.current_batches.push(CurrentBatch { id: id.to_string(), data: BatchData::History { label: tags.label, @@ -576,8 +595,50 @@ impl IrcActor { }, }); } + Some(BatchSubCommand::CUSTOM(c)) if c.as_str() == "DRAFT/MULTILINE" => { + let source = message.source_nickname().unwrap(); + let target = message.response_target().unwrap(); + let msgid = tags.msgid_with_fallback(&["MULTILINE", source, target]); + + let reply = tags + .reply + .as_ref() + .map(|r| { + self.state + .channels + .get(target) + .and_then(|c| c.messages.get(r)) + }) + .map(|m| MessageReference { + text: m.and_then(|m| m.text.clone().map(|t| t.content)), + username: m.map(|m| m.metadata.user.clone()), + }); + + self.current_batches.push(CurrentBatch { + id: id.to_string(), + data: BatchData::Multiline { + target: String::new(), + message: Message { + metadata: MessageMetadata { + msgid, + message_type: MessageType::Privmsg, + server_time: tags.server_time_with_fallback() as f64, + user: source.to_string(), + }, + text: Some(TextMessage { + content: Default::default(), + reactions: Default::default(), + reply, + redacted: false, + edited: false, + relayed_by: tags.relayed_by, + }), + }, + }, + }); + } _ => { - self.current_batch = Some(CurrentBatch { + self.current_batches.push(CurrentBatch { id: id.to_string(), data: BatchData::Unhandled, }); @@ -586,55 +647,106 @@ impl IrcActor { } } else { assert_eq!( - self.current_batch.as_ref().map(|b| b.id.as_str()), + self.current_batches.iter().last().map(|b| b.id.as_str()), Some(&reference[1..]) ); - if let Some(batch) = self.current_batch.take() - && let BatchData::History { - label, - channel: channel_name, - messages, - } = batch.data - { - let channel = self.channel_mut(channel_name.clone()).await; - for message in &messages { - channel - .messages - .insert(message.metadata.msgid.clone(), message.clone()); + if let Some(batch) = self.current_batches.pop() { + match batch.data { + BatchData::History { + label, + channel: channel_name, + messages, + } => { + let channel = self.channel_mut(channel_name.clone()).await; + for message in &messages { + channel + .messages + .insert(message.metadata.msgid.clone(), message.clone()); + } + + let history = History { + channel: channel_name.clone(), + messages, + }; + + let key = if let Some(label) = label { + CommandKey::Label(label) + } else { + CommandKey::History + }; + + self.response_channels + .reply( + &CommandKey::Join(channel_name.clone()), + CommandResponse::Join(channel_name.clone()), + ) + .map_err(|e| { + anyhow!("Failed to reply to JOIN command {e:?}") + })?; + + self.response_channels + .reply(&key, CommandResponse::History(history.clone())) + .unwrap(); + } + + BatchData::Multiline { + target, + message: state_message, + } => { + let source = message.source_nickname().unwrap(); + + if self + .push_batch(target.to_string(), state_message.clone()) + .await + { + return Ok(()); + } + + if let Some(username) = tags.account.clone() { + let user = self.user_mut(source.to_string()).await; + user.username = Some(username); + } + + let channel = self.channel_mut(target.to_string()).await; + channel.messages.insert( + state_message.metadata.msgid.clone(), + state_message.clone(), + ); + + if source == self.state.me.as_ref().unwrap().nickname + && let Err(e) = self.response_channels.reply( + &CommandKey::Privmsg { + target: target.to_string(), + text: state_message + .text + .as_ref() + .unwrap() + .content + .clone(), + }, + CommandResponse::Privmsg(Box::new(state_message.clone())), + ) + { + error!("Failed to reply to PRIVMSG command {e:?}"); + } + + self.on_event(ServerEvent::Privmsg { + channel: target.to_string(), + message: state_message, + }) + .await?; + } + BatchData::Unhandled => (), } - - let history = History { - channel: channel_name.clone(), - messages, - }; - - let key = if let Some(label) = label { - CommandKey::Label(label) - } else { - CommandKey::History - }; - - self.response_channels - .reply( - &CommandKey::Join(channel_name.clone()), - CommandResponse::Join(channel_name.clone()), - ) - .map_err(|e| anyhow!("Failed to reply to JOIN command {e:?}"))?; - - self.response_channels - .reply(&key, CommandResponse::History(history.clone())) - .unwrap(); } - - self.current_batch = None; } } ChannelMODE(ref channel_name, ref mode) => { let target = message.response_target().unwrap(); let source = message.source_nickname().unwrap(); - if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { + if !self.current_batches.iter().any(|b| b.is_chathistory()) { dbg!(target, source, channel_name, mode); } } @@ -642,7 +754,7 @@ impl IrcActor { let target = message.response_target().unwrap(); let source = message.source_nickname().unwrap(); - if self.current_batch.as_ref().map(|b| b.is_chathistory()) != Some(true) { + if !self.current_batches.iter().any(|b| b.is_chathistory()) { dbg!(target, source, channel_name, text); } } @@ -741,7 +853,7 @@ impl IrcActor { } pub async fn push_batch(&mut self, target: String, state_message: Message) -> bool { - if let Some(batch) = self.current_batch.as_mut() + if let Some(batch) = self.current_batches.iter_mut().find(|b| b.is_chathistory()) && let BatchData::History { channel, messages, .. } = &mut batch.data From fe7ef910e7f07de2bdf451c23d617b570d1da483 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 11 Aug 2026 19:26:04 +0200 Subject: [PATCH 25/28] Make IrcChannel::state() return optional --- packages/core/core-wasm/src/lib.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/core/core-wasm/src/lib.rs b/packages/core/core-wasm/src/lib.rs index 981f988..64a160d 100644 --- a/packages/core/core-wasm/src/lib.rs +++ b/packages/core/core-wasm/src/lib.rs @@ -352,7 +352,7 @@ pub struct IrcChannel { #[wasm_bindgen] impl IrcChannel { #[wasm_bindgen] - pub async fn state(&mut self) -> Result { + pub async fn state(&mut self) -> Result, OrbitError> { let (tx, rx) = oneshot::channel(); self.address .send(ActorMessage { @@ -367,7 +367,7 @@ impl IrcChannel { unreachable!("expected state, got: {:?}", resp); }; - Ok(channel.unwrap().into()) + Ok((*channel).map(Into::into)) } #[wasm_bindgen] From fb4d59dfae4bb1dca5aaaea27155c3210ad1ba22 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 11 Aug 2026 19:29:03 +0200 Subject: [PATCH 26/28] Adjust MessageReference doc comment --- packages/core/core-shared/src/state.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/core/core-shared/src/state.rs b/packages/core/core-shared/src/state.rs index abb0753..27aee8c 100644 --- a/packages/core/core-shared/src/state.rs +++ b/packages/core/core-shared/src/state.rs @@ -550,9 +550,9 @@ impl Eq for MessageMetadata {} #[cfg_attr(feature = "web", derive(Tsify))] #[cfg_attr(feature = "web", wasm_bindgen(getter_with_clone, inspectable))] pub struct MessageReference { - /// Empty if message wasn't found or if reply wasn't to a text message + /// Unset if message wasn't found or if reply wasn't to a text message pub text: Option, - /// Empty if message wasn't found + /// Unset if message wasn't found pub username: Option, } From 70a68dca221ae09797446b22fb7fbbadfe04d847 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 11 Aug 2026 19:41:50 +0200 Subject: [PATCH 27/28] Add assert for unfulfilled history requests --- packages/core/core-shared/src/actor.rs | 35 ++++++++++++++++++-------- 1 file changed, 24 insertions(+), 11 deletions(-) diff --git a/packages/core/core-shared/src/actor.rs b/packages/core/core-shared/src/actor.rs index d6370ec..0f86893 100644 --- a/packages/core/core-shared/src/actor.rs +++ b/packages/core/core-shared/src/actor.rs @@ -232,7 +232,7 @@ pub struct IrcActor { disconnect_handlers: Vec>, current_batches: Vec, - requested_history_batches: Vec, + requested_history_batches: Vec<(RequestedHistory, Instant)>, sasl_state: SaslState, rng: SmallRng, } @@ -314,6 +314,13 @@ impl IrcActor { _ = timeout => { self.response_channels.check_timeouts(); + + assert!( + self.requested_history_batches + .iter() + .all(|(_, creation)| creation.elapsed() < Duration::from_secs(5)) + ); + timeout = create_timeout(); } } @@ -388,10 +395,13 @@ impl IrcActor { None }; - self.requested_history_batches.push(RequestedHistory { - channel: channel_name.clone(), - label: label.clone(), - }); + self.requested_history_batches.push(( + RequestedHistory { + channel: channel_name.clone(), + label: label.clone(), + }, + Instant::now(), + )); self.history_latest(channel_name.clone(), None, 5, label) .await .context("Failed to request latest history")?; @@ -582,9 +592,9 @@ impl IrcActor { let idx = self .requested_history_batches .iter() - .position(|b| b.label == tags.label) + .position(|b| b.0.label == tags.label) .expect("Chat history was requested"); - let channel = self.requested_history_batches.remove(idx).channel; + let channel = self.requested_history_batches.remove(idx).0.channel; self.current_batches.push(CurrentBatch { id: id.to_string(), @@ -1123,10 +1133,13 @@ impl IrcActor { None }; - self.requested_history_batches.push(RequestedHistory { - channel: channel.clone(), - label: label.clone(), - }); + self.requested_history_batches.push(( + RequestedHistory { + channel: channel.clone(), + label: label.clone(), + }, + Instant::now(), + )); self.history_before(channel, format!("msgid={before_msgid}"), 5, label) .await .context("Failed to send history before")?; From 75d96bb596f883506104e27e6fe21db529e90a13 Mon Sep 17 00:00:00 2001 From: Jokler Date: Tue, 11 Aug 2026 21:03:02 +0200 Subject: [PATCH 28/28] Fix TS error --- packages/app/src/stores/irc.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/app/src/stores/irc.ts b/packages/app/src/stores/irc.ts index b6bd2a1..7a4634e 100644 --- a/packages/app/src/stores/irc.ts +++ b/packages/app/src/stores/irc.ts @@ -74,7 +74,7 @@ export const useIrcStore = defineStore("irc", () => { console.log("Signed in") // Set initial channel messages - const channelState = await serverChannel.value.state() + const channelState = (await serverChannel.value.state())! const existingServer = serverMessages.get(state.id) ?? new Map() const existingChannel = existingServer.get(channelState.metadata.name) ?? []