use crate::config::Config;
use crate::klaczdb::KlaczDB;
use crate::tools::{ToStringExt, membership_status, room_name};
use std::ops::{Add, Deref};
use std::sync::LazyLock;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use anyhow::bail;
use futures::Future;
use prometheus::{IntCounterVec, opts, register_int_counter_vec};
use tracing::{debug, error, info, trace, warn};
use matrix_sdk::event_handler::{Ctx, EventHandlerHandle};
use matrix_sdk::ruma::OwnedUserId;
use matrix_sdk::ruma::events::room::message::{
MessageType, OriginalSyncRoomMessageEvent, RoomMessageEventContent,
};
use matrix_sdk::{Client, Room};
use tokio::sync::mpsc;
use tokio::task::AbortHandle;
use askama::Template;
use mlua::Lua;
pub static MODULE_EVENTS: LazyLock<IntCounterVec> = LazyLock::new(|| {
register_int_counter_vec!(
opts!(
"module_event_counts",
"Number of events a module has consumed"
),
&["module"]
)
.unwrap()
});
pub static MODULE_ACL_REJECTS: LazyLock<IntCounterVec> = LazyLock::new(|| {
register_int_counter_vec!(
opts!(
"module_acl_failures",
"Number acl checks failed, grouped by module"
),
&["module"]
)
.unwrap()
});
pub static MODULE_CHANNEL_FULL: LazyLock<IntCounterVec> = LazyLock::new(|| {
register_int_counter_vec!(
opts!(
"module_channel_full",
"Number of events a module did not consume due to event channel being full"
),
&["module"]
)
.unwrap()
});
#[derive(Clone)]
pub struct ConsumerEvent {
pub klacz_level: i64,
pub ev: OriginalSyncRoomMessageEvent,
pub sender: OwnedUserId,
pub room: Room,
pub keyword: String,
pub args: Option<String>,
pub lua: Lua,
pub klacz: KlaczDB,
}
#[derive(Clone, Debug)]
pub struct ModuleInfo {
pub name: String,
pub help: String,
pub acl: Vec<Acl>,
pub trigger: TriggerType,
pub channel: mpsc::Sender<ConsumerEvent>,
pub error_prefix: Option<String>,
}
impl ModuleInfo {
pub fn new<C, Fut>(
name: &str,
help: &str,
acl: Vec<Acl>,
trigger: TriggerType,
error_prefix: Option<&str>,
config: C,
processor: impl Fn(ConsumerEvent, C) -> Fut + Send + 'static,
) -> Self
where
C: Clone + Send + Sync + 'static,
Fut: Future<Output = anyhow::Result<()>> + Send + 'static,
{
let owned_error_prefix = error_prefix.map(str::to_owned);
let (tx, rx) = mpsc::channel(1);
Self::spawn_inner(name, owned_error_prefix.clone(), rx, config, processor);
Self {
name: name.to_owned(),
help: help.to_owned(),
acl,
trigger,
channel: tx,
error_prefix: owned_error_prefix,
}
}
fn spawn_inner<C, Fut>(
name: &str,
error_prefix: Option<String>,
rx: mpsc::Receiver<ConsumerEvent>,
config: C,
processor: impl Fn(ConsumerEvent, C) -> Fut + Send + 'static,
) where
C: Clone + Send + Sync + 'static,
Fut: Future<Output = anyhow::Result<()>> + Send + 'static,
{
tokio::task::spawn(Self::consumer(
rx,
config,
error_prefix,
processor,
name.to_owned(),
));
}
pub fn spawn<C, Fut>(
&self,
rx: mpsc::Receiver<ConsumerEvent>,
config: C,
processor: impl Fn(ConsumerEvent, C) -> Fut + Send + 'static,
) where
C: Clone + Send + Sync + 'static,
Fut: Future<Output = anyhow::Result<()>> + Send + 'static,
{
tokio::task::spawn(Self::consumer(
rx,
config,
self.error_prefix.clone(),
processor,
self.name.clone(),
));
}
pub async fn consumer<C, Fut>(
mut rx: mpsc::Receiver<ConsumerEvent>,
config: C,
error_prefix: Option<String>,
processor: impl Fn(ConsumerEvent, C) -> Fut,
name: String,
) -> anyhow::Result<()>
where
C: Clone + Send + Sync,
Fut: Future<Output = anyhow::Result<()>>,
{
loop {
let Some(event) = rx.recv().await else {
warn!("{name} channel closed");
bail!("channel closed");
};
if let Err(e) = processor(event.clone(), config.clone()).await {
error!("error processing event: {e}");
if let Some(ref prefix) = error_prefix {
if let Err(ee) = event
.room
.send(RoomMessageEventContent::text_plain(format!(
"{prefix}: {e}"
)))
.await
{
error!("error when sending event response: {ee}");
}
};
}
}
}
}
#[derive(Clone)]
pub struct PassThroughModuleInfo(pub ModuleInfo);
pub type CatchallDecider = fn(
klaczlevel: i64,
sender: OwnedUserId,
room: &Room,
content: &RoomMessageEventContent,
config: &Config,
) -> anyhow::Result<Consumption>;
#[derive(Clone, Debug)]
pub enum TriggerType {
Keyword(Vec<String>),
Catchall(CatchallDecider),
}
#[derive(Clone, Debug)]
pub enum Acl {
ActiveHswawMember,
MaybeInactiveHswawMember,
SpecificUsers(Vec<String>),
Room(Vec<String>),
KlaczLevel(i64),
Homeserver(Vec<String>),
}
#[derive(PartialEq, Eq, PartialOrd, Ord, Clone, Debug)]
pub enum Consumption {
Reject,
Inclusive,
Passthrough,
Exclusive,
}
#[derive(Clone, Debug, Template)]
#[template(
path = "matrix/help-worker.html",
blocks = ["formatted", "plain"],
)]
pub struct WorkerInfo {
name: String,
help: String,
keyword: String,
helper_module: ModuleInfo,
handle: AbortHandle,
}
impl WorkerInfo {
pub fn new<C, Fut>(
name: &str,
help: &str,
keyword: &str,
mx: Client,
config: C,
worker: impl Fn(Client, C) -> Fut + Send + 'static,
) -> Self
where
C: Clone + Send + Sync + 'static,
Fut: Future<Output = anyhow::Result<()>> + Send + 'static,
{
let handle = tokio::task::spawn(worker(mx, config)).abort_handle();
let (tx, rx) = mpsc::channel(1);
let helper_module = ModuleInfo {
name: name.to_owned(),
help: help.to_owned(),
acl: vec![],
trigger: TriggerType::Keyword(vec![keyword.s()]),
channel: tx,
error_prefix: None,
};
tokio::task::spawn(Self::status_consumer(rx, handle.clone(), name.to_owned()));
Self {
name: name.to_owned(),
help: help.to_owned(),
keyword: keyword.to_owned(),
helper_module,
handle,
}
}
async fn status_consumer(
mut rx: mpsc::Receiver<ConsumerEvent>,
worker_handle: AbortHandle,
name: String,
) -> anyhow::Result<()> {
loop {
let Some(event) = rx.recv().await else {
warn!("{name} channel closed");
worker_handle.abort();
bail!("channel closed");
};
if let Err(e) = event
.room
.send(RoomMessageEventContent::text_plain(format!(
"worker running: {}",
!worker_handle.is_finished()
)))
.await
{
error!("error sending worker status response: {e}");
};
}
}
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
pub async fn dispatcher(
ev: OriginalSyncRoomMessageEvent,
room: Room,
config: Ctx<Config>,
modules: Ctx<Vec<ModuleInfo>>,
passthrough_modules: Ctx<Vec<PassThroughModuleInfo>>,
workers: Ctx<Vec<WorkerInfo>>,
klacz: Ctx<KlaczDB>,
lua: Ctx<Lua>,
) {
use Consumption::{Exclusive, Inclusive, Passthrough, Reject};
use TriggerType::{Catchall, Keyword};
let Some(ev_ts) = ev.origin_server_ts.to_system_time() else {
error!("event timestamp couldn't get parsed to system time");
return;
};
if ev_ts.add(Duration::from_secs(10)) < SystemTime::now() {
debug!("received too old event: {ev_ts:?}");
return;
};
let sender: OwnedUserId = ev.sender.clone();
if config.user_id() == sender || config.ignored().contains(&sender.to_string()) {
return;
}
match ev.content.msgtype {
MessageType::Text(_) | MessageType::Notice(_) => (),
_ => return,
}
let text = ev.content.body();
trace!("new dispatcher: getting klacz permission level");
let klacz_level = match klacz.get_level(&room, &sender).await {
Ok(level) => level,
Err(e) => {
error!("error getting klacz permission level: {e}");
0
}
};
let mut args = text.trim_start().splitn(2, [' ', 'Â ', '\t']);
let first = args.next();
let mut prefixes_all: Vec<String> = config.prefixes();
if let Some(hash) = config.prefixes_restricted() {
prefixes_all.extend(hash.keys().map(std::borrow::ToOwned::to_owned));
};
let mut prefix_selected: Option<String> = None;
trace!("cursed prefix matching");
let (keyword, remainder): (String, Option<String>) = {
first.map_or_else(
|| (String::new(), None),
#[allow(clippy::cognitive_complexity)]
|word| {
trace!("first word exists: {word}");
let (mut kw_candidate, mut remainder_candidate) = (String::new(), None);
for prefix in prefixes_all {
trace!("trying prefix: {prefix}");
match prefix.len() {
1 => match word.strip_prefix(prefix.as_str()) {
None => continue,
Some(w) => {
kw_candidate = w.to_string();
remainder_candidate =
args.next().map(std::string::ToString::to_string);
trace!("selected prefix: {prefix}");
prefix_selected = Some(prefix);
break;
}
},
2.. => {
if word == prefix {
if let Some(shifted_text) = args.next() {
let mut shifted_args =
shifted_text.trim_start().splitn(2, [' ', 'Â ', '\t']);
if let Some(second) = shifted_args.next() {
kw_candidate = second.to_string();
remainder_candidate = shifted_args
.next()
.map(std::string::ToString::to_string);
};
};
trace!("selected prefix: {prefix}");
prefix_selected = Some(prefix);
break;
};
}
0 => continue,
};
}
(kw_candidate, remainder_candidate)
},
)
};
let consumer_event = ConsumerEvent {
klacz_level,
ev: ev.clone(),
sender: sender.clone(),
room: room.clone(),
keyword: keyword.clone(),
args: remainder,
lua: lua.deref().clone(),
klacz: klacz.deref().clone(),
};
let mut run_modules: Vec<(Consumption, ModuleInfo)> = vec![];
let mut consumption = Inclusive;
trace!("figuring out event consumption priority");
for module in modules
.iter()
.chain(workers.iter().map(|w| &w.helper_module))
{
trace!("considering module: {}", module.name);
if module.channel.is_closed() {
debug!("failed module, skipping");
continue;
};
let module_consumption: Consumption = match module.trigger {
Keyword(ref keywords) => {
if !keyword.is_empty() && keywords.contains(&keyword) {
Exclusive
} else {
trace!("\"{keyword}\" doesn't match any keyword: {keywords:?}");
Reject
}
}
Catchall(fun) => match fun(
klacz_level,
sender.clone(),
&room,
&ev.content,
&config.clone(),
) {
Err(e) => {
error!("{} decider returned error: {e}", module.name);
continue;
}
Ok(Reject) => continue,
Ok(c) => c,
},
};
trace!("checking consumption: {module_consumption:?}");
match module_consumption {
Exclusive => {
run_modules.truncate(0);
run_modules.push((module_consumption.clone(), module.clone()));
consumption = module_consumption;
break;
}
Passthrough => {
run_modules.retain(|x| x.0 == Passthrough);
run_modules.push((module_consumption.clone(), module.clone()));
consumption = module_consumption;
}
Inclusive => {
if consumption > module_consumption {
continue;
};
run_modules.push((module_consumption, module.clone()));
}
Reject => continue,
};
}
if consumption == Consumption::Exclusive {
if let Some(prefix) = prefix_selected {
if let Some(map) = config.prefixes_restricted() {
if let Some(list) = map.get(&prefix) {
for (_, module) in &run_modules {
if !list.contains(&module.name) {
return;
}
}
}
}
}
}
trace!("dispatching event to modules");
for (_, module) in run_modules {
dispatch_module(
config.clone(),
true,
&module,
klacz_level,
sender.clone(),
room.clone(),
consumer_event.clone(),
)
.await;
}
match consumption {
Consumption::Inclusive | Consumption::Passthrough => {
trace!("dispatching event to passthrough modules");
let mut run_passthrough_modules: Vec<ModuleInfo> = vec![];
for module in passthrough_modules.iter() {
match module.0.trigger {
Keyword(_) => {
error!("can't have keyword modules in passthrough!");
continue;
}
Catchall(fun) => match fun(
klacz_level,
sender.clone(),
&room,
&ev.content,
&config.clone(),
) {
Err(e) => {
error!("{} decider returned error: {e}", module.0.name);
continue;
}
Ok(Reject) => continue,
Ok(_) => run_passthrough_modules.push(module.0.clone()),
},
};
}
for module in run_passthrough_modules {
if module.channel.is_closed() {
debug!("failed module, skipping");
continue;
};
dispatch_module(
config.clone(),
false,
&module,
klacz_level,
sender.clone(),
room.clone(),
consumer_event.clone(),
)
.await;
}
}
_ => {
debug!("skipping passthrough modules");
}
};
}
#[allow(clippy::too_many_lines)]
pub async fn dispatch_module(
config: Config,
general: bool,
module: &ModuleInfo,
klacz_level: i64,
sender: OwnedUserId,
room: Room,
consumer_event: ConsumerEvent,
) {
use crate::tools::MembershipStatus::{Active, Inactive, NotAMember, Stoned};
use Acl::{
ActiveHswawMember, Homeserver, KlaczLevel, MaybeInactiveHswawMember, Room, SpecificUsers,
};
trace!("dispatching module: {}", module.name);
let mut failed = false;
for acl in &module.acl {
trace!("checking acl: {acl:#?}");
match acl {
KlaczLevel(required) => {
trace!("required: {required}, current: {klacz_level}");
if required > &klacz_level {
failed = true;
}
}
Homeserver(homeservers) => {
if !homeservers.contains(&sender.clone().server_name().to_string()) {
failed = true;
}
}
Room(rooms) => {
let name = room_name(&room);
if !rooms.contains(&name) {
failed = true;
}
}
SpecificUsers(users) => {
if !users.contains(&sender.to_string()) {
failed = true;
}
}
ActiveHswawMember => {
match membership_status(config.capacifier_token(), sender.clone()).await {
Err(e) => {
error!("checking membership for {sender} failed: {e}");
failed = true;
}
Ok(status) => match status {
Inactive | Stoned | NotAMember => failed = true,
Active(_) => (),
},
}
}
MaybeInactiveHswawMember => {
match membership_status(config.capacifier_token(), sender.clone()).await {
Err(e) => {
error!("checking membership for {sender} failed: {e}");
failed = true;
}
Ok(status) => match status {
Stoned | NotAMember => failed = true,
Inactive | Active(_) => (),
},
}
}
};
if failed {
break;
}
}
if failed {
MODULE_ACL_REJECTS.with_label_values(&[&module.name]).inc();
if general {
let mut response = "busy figuring out why time behaves weirdly";
let options = config.acl_deny();
if let Ok(now) = SystemTime::now().duration_since(UNIX_EPOCH) {
let milis = now.as_millis();
let chosen_idx: usize = milis as usize % config.acl_deny().len();
if let Some(option) = options.get(chosen_idx) {
response = option;
};
};
if let Err(e) = room
.send(RoomMessageEventContent::text_plain(response))
.await
{
error!("sending acl failure response failed: {e}");
}
};
return;
};
if config.modules_disabled().contains(&module.name)
|| config.modules_fenced().contains(&module.name)
{
trace!("module disabled: {}", module.name);
return;
}
trace!("attempting to reserve channel space");
let reservation = match module.channel.clone().try_reserve_owned() {
Ok(r) => r,
Err(e) => {
MODULE_CHANNEL_FULL.with_label_values(&[&module.name]).inc();
error!("module {} channel can't accept message: {e}", module.name);
return;
}
};
MODULE_EVENTS.with_label_values(&[&module.name]).inc();
trace!("sending event");
reservation.send(consumer_event);
}
#[allow(
clippy::cognitive_complexity,
reason = "Just a few loops, heurestics seem wrong here"
)]
pub fn init_modules(
mx: &Client,
config: &Config,
reload_tx: mpsc::Sender<Room>,
) -> EventHandlerHandle {
let klacz = KlaczDB { handle: "main" };
let mut modules: Vec<ModuleInfo> = vec![];
let mut passthrough_modules: Vec<PassThroughModuleInfo> = vec![];
let mut workers: Vec<WorkerInfo> = vec![];
if let Err(e) = crate::notmun::module_starter(mx, config) {
error!("failed initializing notmun: {e}");
} else {
info!("initialized notmun");
};
for starter in [
crate::klaczdb::starter,
crate::spaceapi::starter,
crate::db::starter,
crate::inviter::starter,
crate::kasownik::starter,
crate::wolfram::starter,
crate::sage::starter,
crate::alerts::starter,
crate::autojoiner::starter,
crate::forgejo::starter,
] {
match starter(mx, config) {
Err(e) => error!("module initialization failed fatally: {e}"),
Ok(m) => modules.extend(m),
};
}
for starter in [crate::kasownik::passthrough, crate::notmun::passthrough] {
match starter(mx, config) {
Err(e) => error!("module initialization failed fatally: {e}"),
Ok(m) => passthrough_modules.extend(m),
};
}
for starter in [
crate::webterface::workers,
crate::spaceapi::workers,
crate::forgejo::workers,
crate::gerrit::workers,
] {
match starter(mx, config) {
Err(e) => error!("module initialization failed fatally: {e}"),
Ok(m) => workers.extend(m),
};
}
modules.retain(|x| !config.modules_fenced().contains(&x.name));
passthrough_modules.retain(|x| !config.modules_fenced().contains(&x.0.name));
modules.extend(core_starter(
config,
reload_tx,
&modules,
&passthrough_modules,
workers.clone(),
));
mx.add_event_handler_context(klacz);
mx.add_event_handler_context(config.clone());
mx.add_event_handler_context(modules);
mx.add_event_handler_context(passthrough_modules);
mx.add_event_handler_context(workers);
mx.add_event_handler(dispatcher)
}
pub fn core_starter(
config: &Config,
reload_ev_tx: mpsc::Sender<Room>,
registered_modules: &[ModuleInfo],
registered_passthrough_modules: &[PassThroughModuleInfo],
registered_workers: Vec<WorkerInfo>,
) -> Vec<ModuleInfo> {
info!("registering modules");
let mut modules: Vec<ModuleInfo> = vec![];
let (help_tx, help_rx) = mpsc::channel::<ConsumerEvent>(1);
let help = ModuleInfo {
name: "help".s(),
help: "get help about the bot or its basic functions".s(),
acl: vec![],
trigger: TriggerType::Keyword(vec!["help".s(), "status".s()]),
channel: help_tx,
error_prefix: None,
};
modules.push(help);
let (list_tx, list_rx) = mpsc::channel::<ConsumerEvent>(1);
let list = ModuleInfo {
name: "list".s(),
help: "get the list of currently registered modules".s(),
acl: vec![],
trigger: TriggerType::Keyword(vec!["list".s(), "list-functions".s()]),
channel: list_tx,
error_prefix: None,
};
modules.push(list);
let (reload_tx, reload_rx) = mpsc::channel::<ConsumerEvent>(1);
let reload = ModuleInfo {
name: "reload".s(),
help: "reload bot configuration and modules".s(),
acl: vec![Acl::SpecificUsers(config.admins())],
trigger: TriggerType::Keyword(vec!["reload".s()]),
channel: reload_tx,
error_prefix: None,
};
modules.push(reload);
let (shutdown_tx, shutdown_rx) = mpsc::channel::<ConsumerEvent>(1);
let shutdown = ModuleInfo {
name: "shutdown".s(),
help: "makes the bot process exit, literally".s(),
acl: vec![Acl::SpecificUsers(config.admins())],
trigger: TriggerType::Keyword(vec!["shutdown".s(), "die".s(), "exit".s()]),
channel: shutdown_tx,
error_prefix: None,
};
modules.push(shutdown);
let mod_manager = ModuleInfo::new(
"mod_manager",
"fences off/disables/unfences/enables modules",
vec![Acl::SpecificUsers(config.admins())],
TriggerType::Keyword(vec![
"enable".s(),
"disable".s(),
"fence".s(),
"unfence".s(),
"disabled".s(),
"fenced".s(),
]),
Some("action failed"),
config.clone(),
mod_manager,
);
modules.push(mod_manager);
let weak_modules: Vec<WeakModuleInfo> = registered_modules
.iter()
.chain(&modules)
.map(std::convert::Into::into)
.collect();
let weak_passthrough: Vec<WeakModuleInfo> = registered_passthrough_modules
.iter()
.map(std::convert::Into::into)
.collect();
tokio::task::spawn(help_consumer(
help_rx,
config.clone(),
weak_modules.clone(),
weak_passthrough.clone(),
registered_workers.clone(),
));
tokio::task::spawn(list_consumer(
list_rx,
weak_modules,
weak_passthrough,
registered_workers,
));
tokio::task::spawn(reload_consumer(reload_rx, reload_ev_tx));
tokio::task::spawn(shutdown_consumer(shutdown_rx));
modules
}
#[derive(Clone, Debug, Template)]
#[template(
path = "matrix/help-module.html",
blocks = ["formatted", "plain"],
)]
pub struct WeakModuleInfo {
pub name: String,
pub help: String,
pub trigger: TriggerType,
pub channel: mpsc::WeakSender<ConsumerEvent>,
}
impl From<&ModuleInfo> for WeakModuleInfo {
fn from(m: &ModuleInfo) -> Self {
Self {
name: m.name.clone(),
help: m.help.clone(),
trigger: m.trigger.clone(),
channel: m.channel.downgrade(),
}
}
}
impl From<&PassThroughModuleInfo> for WeakModuleInfo {
fn from(m: &PassThroughModuleInfo) -> Self {
Self {
name: m.0.name.clone(),
help: m.0.help.clone(),
trigger: m.0.trigger.clone(),
channel: m.0.channel.downgrade(),
}
}
}
pub async fn shutdown_consumer(mut rx: mpsc::Receiver<ConsumerEvent>) -> anyhow::Result<()> {
if rx.recv().await.is_none() {
warn!("shutdown channel closed");
bail!("channel closed");
};
info!("received process exit request");
std::process::exit(0);
}
async fn help_consumer(
mut rx: mpsc::Receiver<ConsumerEvent>,
config: Config,
modules: Vec<WeakModuleInfo>,
passthrough_modules: Vec<WeakModuleInfo>,
workers: Vec<WorkerInfo>,
) -> anyhow::Result<()> {
loop {
let Some(event) = rx.recv().await else {
warn!("help channel closed");
bail!("channel closed");
};
if let Err(e) = help_processor(
event.clone(),
config.clone(),
modules.clone(),
passthrough_modules.clone(),
workers.clone(),
)
.await
{
if let Err(e) = event
.room
.send(RoomMessageEventContent::text_plain(format!(
"error getting help: {e}"
)))
.await
{
error!("error while sending response: {e}");
};
}
}
}
#[derive(Template)]
#[template(
path = "matrix/help-generic.html",
blocks = ["formatted", "plain"],
)]
struct RenderHelp {
config: Config,
modules: Vec<WeakModuleInfo>,
passthrough: Vec<WeakModuleInfo>,
workers: Vec<WorkerInfo>,
source_url: String,
docs_link: String,
matrix_contact: String,
}
impl RenderHelp {
fn failed(&self) -> (usize, usize, usize) {
(
self.modules
.iter()
.filter(|x| x.channel.upgrade().unwrap().is_closed())
.count(),
self.passthrough
.iter()
.filter(|x| x.channel.upgrade().unwrap().is_closed())
.count(),
self.workers
.iter()
.filter(|x| x.handle.is_finished())
.count(),
)
}
}
pub async fn help_processor(
event: ConsumerEvent,
config: Config,
modules: Vec<WeakModuleInfo>,
passthrough: Vec<WeakModuleInfo>,
workers: Vec<WorkerInfo>,
) -> anyhow::Result<()> {
let generic = RenderHelp {
config,
modules: modules.clone(),
passthrough: passthrough.clone(),
workers: workers.clone(),
source_url: "https://code.hackerspace.pl/ar/notbot".s(),
docs_link: "https://docs.rs/notbot/latest/notbot/".s(),
matrix_contact: "@ar:is-a.cat".s(),
};
let generic_help = RoomMessageEventContent::text_html(
generic.as_plain().render()?,
generic.as_formatted().render()?,
);
let Some(args) = event.args else {
event.room.send(generic_help).await?;
return Ok(());
};
let mut arguments = args.split_whitespace();
let Some(maybe_module_name) = arguments.next() else {
event.room.send(generic_help).await?;
return Ok(());
};
let mut specific_response: Option<RoomMessageEventContent> = None;
for module in modules.iter().chain(&passthrough) {
if module.name == maybe_module_name {
let specific_help = RoomMessageEventContent::text_html(
module.as_plain().render()?,
module.as_formatted().render()?,
);
specific_response = Some(specific_help);
break;
};
}
for module in workers {
if module.name == maybe_module_name {
let specific_help = RoomMessageEventContent::text_html(
module.as_plain().render()?,
module.as_formatted().render()?,
);
specific_response = Some(specific_help);
break;
};
}
if let Some(response) = specific_response {
event.room.send(response).await?;
} else {
event.room.send(generic_help).await?;
};
Ok(())
}
#[derive(Template)]
#[template(
path = "matrix/help-list.html",
blocks = ["formatted", "plain"],
)]
struct RenderList {
modules: Vec<WeakModuleInfo>,
passthrough: Vec<WeakModuleInfo>,
workers: Vec<WorkerInfo>,
}
impl RenderList {
#[allow(clippy::unused_self, reason = "required by templating engine")]
fn list_modules(&self, m: &[WeakModuleInfo]) -> (Vec<String>, bool) {
let mut failed = false;
(
m.iter()
.map(|x| {
let mut s = x.name.clone();
if x.channel.upgrade().unwrap().is_closed() {
s.push('*');
failed = true;
}
s
})
.collect(),
failed,
)
}
fn list_workers(&self) -> (Vec<String>, bool) {
let mut failed = false;
(
self.workers
.iter()
.map(|x| {
let mut s = x.name.clone();
if x.handle.is_finished() {
s.push('*');
failed = true;
};
s
})
.collect(),
failed,
)
}
}
pub async fn list_consumer(
mut rx: mpsc::Receiver<ConsumerEvent>,
modules: Vec<WeakModuleInfo>,
passthrough: Vec<WeakModuleInfo>,
workers: Vec<WorkerInfo>,
) -> anyhow::Result<()> {
loop {
let Some(event) = rx.recv().await else {
warn!("list channel closed");
bail!("channel closed");
};
let render_list = RenderList {
modules: modules.clone(),
passthrough: passthrough.clone(),
workers: workers.clone(),
};
let response = RoomMessageEventContent::text_html(
render_list.as_plain().render()?,
render_list.as_formatted().render()?,
);
if let Err(e) = event.room.send(response).await {
error!("failed sending list response: {e}");
}
}
}
pub async fn reload_consumer(
mut rx: mpsc::Receiver<ConsumerEvent>,
reload_tx: mpsc::Sender<Room>,
) -> anyhow::Result<()> {
loop {
let Some(event) = rx.recv().await else {
warn!("reload channel closed");
bail!("channel closed");
};
let reservation = match reload_tx.clone().try_reserve_owned() {
Ok(r) => r,
Err(e) => {
error!("reloader can't accept trigger: {e}");
continue;
}
};
reservation.send(event.room);
}
}
pub async fn mod_manager(event: ConsumerEvent, config: Config) -> anyhow::Result<()> {
let modname = match event.args {
None => match event.keyword.as_str() {
"fenced" | "disabled" => "".s(),
_ => bail!("no module name provided"),
},
Some(m) => m.trim().s(),
};
match event.keyword.as_str() {
"disable" => config.disable_module(modname.clone()),
"enable" => config.enable_module(&modname),
"fence" => config.fence_module(modname.clone()),
"unfence" => config.unfence_module(&modname),
"disabled" => {
let disabled = config.modules_disabled();
let message = format!("disabled modules: {disabled:?}");
event
.room
.send(RoomMessageEventContent::text_plain(message))
.await?;
return Ok(());
}
"fenced" => {
let fenced = config.modules_fenced();
let message = format!("fenced modules: {fenced:?}");
event
.room
.send(RoomMessageEventContent::text_plain(message))
.await?;
return Ok(());
}
_ => bail!("wtf? wrong keyword passed somehow"),
}?;
let message = format!("module {} successfully {}d", modname, event.keyword);
event
.room
.send(RoomMessageEventContent::text_plain(message))
.await?;
Ok(())
}