5 Commits

7 changed files with 1040 additions and 36 deletions
+1
View File
@@ -1 +1,2 @@
/target /target
cliff.toml
Generated
+902 -7
View File
File diff suppressed because it is too large Load Diff
+9 -3
View File
@@ -1,16 +1,21 @@
[package] [package]
name = "RhythmWorks" name = "RhythmWorks"
version = "0.5.0" version = "0.7.0"
edition = "2024" edition = "2024"
[dependencies] [dependencies]
color-eyre = "0.6.5"
crossterm = "0.29.0"
dialoguer = "0.12.0" dialoguer = "0.12.0"
keyring = { version = "3.6.3", features = ["linux-native", "windows-native"] } keyring = { version = "3.6.3", features = ["linux-native", "windows-native"] }
opentelemetry = { version = "0.27.0", features = ["logs", "metrics", "trace"] } opentelemetry = { version = "0.27.0", features = ["logs", "metrics", "trace"] }
opentelemetry-otlp = { version = "0.27.0", features = ["trace", "metrics", "grpc-tonic", "http-proto", "tls", "reqwest-client", "reqwest-rustls"] } opentelemetry-appender-log = "0.27.0"
opentelemetry-appender-tracing = "0.27.0"
opentelemetry-otlp = { version = "0.27.0", features = ["trace", "metrics", "grpc-tonic", "http-proto", "tls", "reqwest-client", "reqwest-rustls", "logs"] }
opentelemetry-proto = "0.27.0" opentelemetry-proto = "0.27.0"
opentelemetry-semantic-conventions = "0.27.0" opentelemetry-semantic-conventions = "0.27.0"
opentelemetry_sdk = { version = "0.27.0", features = ["rt-tokio", "trace"] } opentelemetry_sdk = { version = "0.27.0", features = ["logs", "rt-tokio", "trace"] }
ratatui = { version = "0.30.0", features = ["all-widgets"] }
reqwest = { version = "0.12.28", features = ["blocking", "json"] } reqwest = { version = "0.12.28", features = ["blocking", "json"] }
semver = "1.0.27" semver = "1.0.27"
serde = "1.0.228" serde = "1.0.228"
@@ -19,6 +24,7 @@ serenity = { version = "0.12.5", features = ["client", "gateway", "voice"] }
songbird = { version = "0.5.0", features = ["builtin-queue", "driver", "serenity"] } songbird = { version = "0.5.0", features = ["builtin-queue", "driver", "serenity"] }
symphonia = { version = "0.5.5", features = ["aac", "alac", "isomp4", "mp3"] } symphonia = { version = "0.5.5", features = ["aac", "alac", "isomp4", "mp3"] }
tokio = "1.48.0" tokio = "1.48.0"
tokio-util = "0.7.18"
tonic = { version = "0.12.3", features = ["tls-roots"] } tonic = { version = "0.12.3", features = ["tls-roots"] }
tracing = "0.1.41" tracing = "0.1.41"
tracing-opentelemetry = "0.32.0" tracing-opentelemetry = "0.32.0"
+1
View File
@@ -0,0 +1 @@
## [0.5.2] - 2026-01-21
+102 -19
View File
@@ -1,4 +1,3 @@
use core::error;
use dialoguer::Select; use dialoguer::Select;
use dialoguer::theme::ColorfulTheme; use dialoguer::theme::ColorfulTheme;
use modules::*; use modules::*;
@@ -19,15 +18,28 @@ use songbird::events::EventHandler as VoiceEventHandler;
use songbird::get; use songbird::get;
use songbird::input::HlsRequest; use songbird::input::HlsRequest;
use songbird::input::Input; use songbird::input::Input;
use tracing::instrument;
use std::env; use std::env;
use std::process::Command; use std::process::Command;
use std::process::Stdio; use std::process::Stdio;
pub mod modules; use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use tracing::level_filters::LevelFilter;
use tracing_subscriber::Layer;
use tracing_subscriber::filter::Targets;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
pub mod modules;
enum BotState {
Stopped,
Running,
}
struct TrackTraceHandler { struct TrackTraceHandler {
otel_ctx: opentelemetry::Context, otel_ctx: opentelemetry::Context,
guild_id: String, guild_id: String,
track_url: String, track_url: String,
hls_url: String,
event_name: &'static str, event_name: &'static str,
} }
@@ -45,7 +57,15 @@ impl VoiceEventHandler for TrackTraceHandler {
.set_attribute(KeyValue::new("Guild_ID", self.guild_id.clone())); .set_attribute(KeyValue::new("Guild_ID", self.guild_id.clone()));
cx.span() cx.span()
.set_attribute(KeyValue::new("Youtube_URL", self.track_url.clone())); .set_attribute(KeyValue::new("Youtube_URL", self.track_url.clone()));
cx.span().add_event(self.event_name, vec![KeyValue::new("track_status", track_status.unwrap().to_string())]); cx.span()
.set_attribute(KeyValue::new("Stream_URL", self.hls_url.clone()));
cx.span().add_event(
self.event_name,
vec![KeyValue::new(
"track_status",
track_status.unwrap().to_string(),
)],
);
}); });
None None
@@ -75,6 +95,7 @@ impl EventHandler for Handler {
let cx = opentelemetry::Context::current_with_span(span); let cx = opentelemetry::Context::current_with_span(span);
let guild_id = match msg.guild_id { let guild_id = match msg.guild_id {
Some(g) => { Some(g) => {
tracing::info!("Received Guild ID");
cx.span().add_event( cx.span().add_event(
"Received Guild ID", "Received Guild ID",
vec![KeyValue::new("Server ID", g.to_string())], vec![KeyValue::new("Server ID", g.to_string())],
@@ -215,7 +236,7 @@ impl EventHandler for Handler {
manager.join(guild_id, channel_id).await.unwrap() manager.join(guild_id, channel_id).await.unwrap()
}; };
let output = match Command::new("./yt-dlp") let output = match Command::new("./yt-dlp")
.args(["-f", "bestaudio", "-g", &yt_url]) .args(["-f", "bestaudio", "-g", "--no-playlist", &yt_url])
.stdout(Stdio::piped()) .stdout(Stdio::piped())
.output() .output()
{ {
@@ -240,6 +261,7 @@ impl EventHandler for Handler {
.expect("Invalid UTF-8") .expect("Invalid UTF-8")
.trim() .trim()
.to_string(); .to_string();
let telemetry_stream_url = stream_url.clone();
let mut call = handler_lock.lock().await; let mut call = handler_lock.lock().await;
let input = Input::from(HlsRequest::new(reqwest::Client::new(), stream_url)); let input = Input::from(HlsRequest::new(reqwest::Client::new(), stream_url));
let track_handle = call.enqueue_input(input).await; let track_handle = call.enqueue_input(input).await;
@@ -265,6 +287,7 @@ impl EventHandler for Handler {
otel_ctx: otel_ctx.clone(), otel_ctx: otel_ctx.clone(),
guild_id: guild_id.to_string().clone(), guild_id: guild_id.to_string().clone(),
track_url: yt_url.clone(), track_url: yt_url.clone(),
hls_url: telemetry_stream_url.clone(),
event_name: "Track Started Playing", event_name: "Track Started Playing",
}, },
); );
@@ -274,6 +297,7 @@ impl EventHandler for Handler {
otel_ctx: otel_ctx.clone(), otel_ctx: otel_ctx.clone(),
guild_id: guild_id.to_string().clone(), guild_id: guild_id.to_string().clone(),
track_url: yt_url.clone(), track_url: yt_url.clone(),
hls_url: telemetry_stream_url.clone(),
event_name: "Track Finished", event_name: "Track Finished",
}, },
); );
@@ -283,6 +307,7 @@ impl EventHandler for Handler {
otel_ctx: otel_ctx.clone(), otel_ctx: otel_ctx.clone(),
guild_id: guild_id.to_string().clone(), guild_id: guild_id.to_string().clone(),
track_url: yt_url.clone(), track_url: yt_url.clone(),
hls_url: telemetry_stream_url.clone(),
event_name: "Playback had an error", event_name: "Playback had an error",
}, },
); );
@@ -337,6 +362,30 @@ impl EventHandler for Handler {
println!("{} is connected!", ready.user.name); println!("{} is connected!", ready.user.name);
} }
} }
async fn run_bot(token: String, shutdown: CancellationToken) {
let intents = GatewayIntents::GUILD_MESSAGES
| GatewayIntents::GUILDS
| GatewayIntents::DIRECT_MESSAGES
| GatewayIntents::MESSAGE_CONTENT
| GatewayIntents::GUILD_VOICE_STATES;
let mut client = Client::builder(&token, intents)
.event_handler(Handler)
.register_songbird()
.await
.expect("Err creating client");
tokio::select! {
result = client.start() => {
if let Err(why) = result {
println!("Client error: {why:?}");
}
}
_ = shutdown.cancelled() => {
println!("Shutting down bot...");
}
}
}
#[tokio::main] #[tokio::main]
async fn main() { async fn main() {
@@ -345,17 +394,33 @@ async fn main() {
This program comes with ABSOLUTELY NO WARRANTY\n This program comes with ABSOLUTELY NO WARRANTY\n
This is free software, and you are welcome to redistribute it under certain conditions." This is free software, and you are welcome to redistribute it under certain conditions."
); );
let tracer_provider = telemetry::init_telemetry(); let (tracer_provider, logger_provider) = telemetry::init_telemetry();
global::set_tracer_provider(tracer_provider.clone()); global::set_tracer_provider(tracer_provider.clone());
let filter = Targets::new()
.with_target("RhythmWorks", LevelFilter::INFO)
.with_default(LevelFilter::OFF);
let otel_log_layer =
opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge::new(&logger_provider)
.with_filter(filter);
tracing_subscriber::registry().with(otel_log_layer).init();
let tracer: global::BoxedTracer = global::tracer("tracer"); let tracer: global::BoxedTracer = global::tracer("tracer");
tracer tracer
.in_span("Checking for Updates", |_cx| updator::update()) .in_span("Checking for Updates", |_cx| updator::update())
.await; .await;
let service_name = "RhythmWorks"; let service_name = "RhythmWorks";
let username = env::var("USER").or_else(|_| env::var("USERNAME")).unwrap(); let username = env::var("USER").or_else(|_| env::var("USERNAME")).unwrap();
let mut bot_task: Option<JoinHandle<()>> = None;
let mut shutdown_token: Option<CancellationToken> = None;
let mut bot_state = BotState::Stopped;
loop { loop {
let token = credentialmanager::get_token(service_name, &username); let token = credentialmanager::get_token(service_name, &username);
let options: Vec<&'static str> = vec!["- Start Bot", "- Change Token", "- Delete Token"]; let options: Vec<&'static str> = vec![
"- Start Bot",
"- Stop Bot",
"- Change Token",
"- Delete Token",
"- Exit",
];
let selection = Select::with_theme(&ColorfulTheme::default()) let selection = Select::with_theme(&ColorfulTheme::default())
.with_prompt("Choose an Option") .with_prompt("Choose an Option")
.default(0) .default(0)
@@ -364,26 +429,44 @@ This is free software, and you are welcome to redistribute it under certain cond
.unwrap(); .unwrap();
match selection { match selection {
0 => { 0 => {
let intents = GatewayIntents::GUILD_MESSAGES if matches!(bot_state, BotState::Running) {
| GatewayIntents::GUILDS println!("[x] Bot is already running.");
| GatewayIntents::DIRECT_MESSAGES continue;
| GatewayIntents::MESSAGE_CONTENT
| GatewayIntents::GUILD_VOICE_STATES;
let mut client = Client::builder(&token, intents)
.event_handler(Handler)
.register_songbird()
.await
.expect("Err creating client");
if let Err(why) = client.start().await {
println!("Client error: {why:?}");
} }
let token = credentialmanager::get_token(service_name, &username);
let shutdown = CancellationToken::new();
let shutdown_clone = shutdown.clone();
let task = tokio::spawn(run_bot(token, shutdown_clone));
bot_task = Some(task);
shutdown_token = Some(shutdown);
bot_state = BotState::Running;
println!("✅ Bot started.");
} }
1 => { 1 => {
credentialmanager::change_token(service_name, &username); if let Some(token) = shutdown_token.take() {
token.cancel();
}
if let Some(task) = bot_task.take() {
let _ = task.await;
}
println!("Bot stopped.");
} }
2 => { 2 => {
credentialmanager::change_token(service_name, &username);
}
3 => {
credentialmanager::remove_token(service_name, &username); credentialmanager::remove_token(service_name, &username);
} }
4 => {
if let Some(token) = shutdown_token {
token.cancel();
}
break;
}
_ => println!("Invalid Option"), _ => println!("Invalid Option"),
} }
} }
+23 -5
View File
@@ -1,9 +1,13 @@
use opentelemetry::KeyValue; use opentelemetry::KeyValue;
use opentelemetry::global;
use opentelemetry_otlp::{WithExportConfig, WithTonicConfig}; use opentelemetry_otlp::{WithExportConfig, WithTonicConfig};
use opentelemetry_sdk::{Resource, runtime}; use opentelemetry_sdk::{Resource, logs, runtime};
use tonic::transport::{Channel, ClientTlsConfig}; use tonic::transport::{Channel, ClientTlsConfig};
pub fn init_telemetry() -> opentelemetry_sdk::trace::TracerProvider { pub fn init_telemetry() -> (
opentelemetry_sdk::trace::TracerProvider,
opentelemetry_sdk::logs::LoggerProvider,
) {
let endpoint = "https://signoz.racooncity.org".to_string(); let endpoint = "https://signoz.racooncity.org".to_string();
let channel = Channel::from_shared(endpoint.clone()) let channel = Channel::from_shared(endpoint.clone())
.unwrap() .unwrap()
@@ -13,14 +17,28 @@ pub fn init_telemetry() -> opentelemetry_sdk::trace::TracerProvider {
let exporter = opentelemetry_otlp::SpanExporter::builder() let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_tonic() .with_tonic()
.with_endpoint(endpoint.clone()) .with_endpoint(endpoint.clone())
.with_channel(channel) .with_channel(channel.clone())
.build() .build()
.expect("Failed to Build Exporter"); .expect("Failed to Build Exporter");
opentelemetry_sdk::trace::TracerProvider::builder() let log_exporter = opentelemetry_otlp::LogExporter::builder()
.with_tonic()
.with_endpoint(endpoint.clone())
.with_channel(channel.clone())
.build()
.expect("Failed to build log exporter");
let logger_provider = logs::LoggerProvider::builder()
.with_resource(Resource::new(vec![KeyValue::new(
"service.name",
"RhythmWorks",
)]))
.with_batch_exporter(log_exporter, runtime::Tokio)
.build();
let tracer_provider = opentelemetry_sdk::trace::TracerProvider::builder()
.with_batch_exporter(exporter, runtime::Tokio) .with_batch_exporter(exporter, runtime::Tokio)
.with_resource(Resource::new(vec![KeyValue::new( .with_resource(Resource::new(vec![KeyValue::new(
"service.name", "service.name",
"RhythmWorks", "RhythmWorks",
)])) )]))
.build() .build();
(tracer_provider, logger_provider)
} }
+2 -2
View File
@@ -10,7 +10,7 @@ struct Release {
pub async fn update() { pub async fn update() {
let body = reqwest::get( let body = reqwest::get(
"https://git.racooncity.org/api/v1/repos/brotoskyj/RhythmWorks/releases/latest", "https://git.racooncity.org/api/v1/repos/mysticmomba/RhythmWorks/releases/latest",
) )
.await .await
.unwrap() .unwrap()
@@ -22,7 +22,7 @@ pub async fn update() {
let current_version = Version::parse(env!("CARGO_PKG_VERSION")).unwrap(); let current_version = Version::parse(env!("CARGO_PKG_VERSION")).unwrap();
if latest_version > current_version { if latest_version > current_version {
println!("Update Available! Please download the latest version"); println!("Update Available! Please download the latest version");
println!("URL: https://git.racooncity.org/brotoskyj/RhythmWorks/releases"); println!("URL: https://git.racooncity.org/mysticmomba/RhythmWorks/releases");
thread::sleep(std::time::Duration::from_secs(10)); thread::sleep(std::time::Duration::from_secs(10));
std::process::exit(1); std::process::exit(1);
} }