Compare commits

...
Author SHA1 Message Date
tdgao cc48fab517 feat: implement members breakdown and filter 2026-06-15 14:08:20 -07:00
2 changed files with 169 additions and 6 deletions
@@ -4,14 +4,18 @@ use serde::{Deserialize, Serialize};
use sqlx::Row;
use crate::{
database::{PgPool, models::DBProjectId},
database::{
PgPool,
models::{DBProjectId, DBUserId},
},
models::ids::ProjectId,
routes::ApiError,
util::error::Context,
};
use ariadne::ids::UserId;
use super::super::{TimeSlice, add_to_time_slice};
use super::{AnalyticsData, ProjectAnalytics, ProjectMetrics};
use super::{AnalyticsData, Metrics, ProjectAnalytics, ProjectMetrics};
/// Fields for [`super::ReturnMetrics::project_revenue`].
#[derive(
@@ -21,15 +25,24 @@ use super::{AnalyticsData, ProjectAnalytics, ProjectMetrics};
pub enum ProjectRevenueField {
/// Project ID.
ProjectId,
/// User ID.
UserId,
}
/// Filters for [`super::ReturnMetrics::project_revenue`].
#[derive(Debug, Clone, Default, Serialize, Deserialize, utoipa::ToSchema)]
pub struct ProjectRevenueFilters {}
pub struct ProjectRevenueFilters {
/// User IDs to include.
#[serde(default)]
pub user_id: Vec<UserId>,
}
/// [`super::ReturnMetrics::project_revenue`].
#[derive(Debug, Clone, Default, Serialize, Deserialize, utoipa::ToSchema)]
pub struct ProjectRevenue {
/// [`ProjectRevenueField::UserId`].
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) user_id: Option<UserId>,
/// Total revenue for this bucket.
pub(crate) revenue: Decimal,
}
@@ -40,7 +53,15 @@ pub(crate) async fn fetch(
req: &super::super::GetRequest,
num_time_slices: usize,
project_id_values: &[i64],
metrics: &Metrics<ProjectRevenueField, ProjectRevenueFilters>,
) -> Result<(), ApiError> {
let use_user_id = metrics.bucket_by.contains(&ProjectRevenueField::UserId);
let filter_user_ids = metrics
.filter_by
.user_id
.iter()
.map(|id| DBUserId::from(*id).0)
.collect::<Vec<_>>();
let mut rows = sqlx::query(
"SELECT
WIDTH_BUCKET(
@@ -50,6 +71,7 @@ pub(crate) async fn fetch(
$3::integer
) AS bucket,
mod_id,
CASE WHEN $5 THEN user_id ELSE 0 END AS member_user_id,
SUM(amount) amount_sum
FROM payouts_values
WHERE
@@ -59,12 +81,15 @@ pub(crate) async fn fetch(
AND payouts_values.mod_id = ANY($4)
AND created >= $1
AND created < $2
GROUP BY bucket, mod_id",
AND (cardinality($6::bigint[]) = 0 OR user_id = ANY($6))
GROUP BY bucket, mod_id, member_user_id",
)
.bind(req.time_range.start)
.bind(req.time_range.end)
.bind(num_time_slices as i64)
.bind(project_id_values)
.bind(use_user_id)
.bind(&filter_user_ids)
.fetch(pool);
while let Some(row) = rows.next().await.transpose()? {
let bucket = row
@@ -77,6 +102,7 @@ pub(crate) async fn fetch(
})?;
let mod_id = row.try_get::<Option<i64>, _>("mod_id")?;
let user_id = row.try_get::<Option<i64>, _>("member_user_id")?;
let amount_sum = row.try_get::<Option<Decimal>, _>("amount_sum")?;
if let Some(source_project) =
mod_id.map(DBProjectId).map(ProjectId::from)
@@ -88,6 +114,10 @@ pub(crate) async fn fetch(
AnalyticsData::Project(ProjectAnalytics {
source_project,
metrics: ProjectMetrics::Revenue(ProjectRevenue {
user_id: user_id
.filter(|id| *id != 0)
.map(DBUserId)
.map(UserId::from),
revenue,
}),
}),
@@ -18,6 +18,7 @@ use std::{
use crate::database::PgPool;
use actix_web::{HttpRequest, post, web};
use ariadne::ids::UserId;
use chrono::{DateTime, TimeDelta, Utc};
use eyre::eyre;
use serde::{Deserialize, Serialize};
@@ -43,7 +44,7 @@ use crate::{
projects::ProjectStatus,
teams::ProjectPermissions,
threads::MessageBody,
v3::{analytics::DownloadReason, projects::Project},
v3::{analytics::DownloadReason, projects::Project, users::User},
},
queue::session::AuthQueue,
routes::ApiError,
@@ -131,6 +132,9 @@ pub struct GetResponse {
/// Project metadata for projects referenced in the response metrics.
#[serde(default)]
pub projects: HashMap<ProjectId, Project>,
/// User metadata for users referenced in the response metrics.
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub users: HashMap<UserId, User>,
/// List of events associated with projects that were requested.
pub project_events: Vec<ProjectAnalyticsEvent>,
}
@@ -355,7 +359,7 @@ pub async fn fetch_analytics(
.await?;
}
if req.return_metrics.project_revenue.is_some() {
if let Some(metrics) = &req.return_metrics.project_revenue {
if !scopes.contains(Scopes::PAYOUTS_READ) {
return Err(AuthenticationError::InvalidCredentials.into());
}
@@ -366,6 +370,7 @@ pub async fn fetch_analytics(
&req,
num_time_slices,
&project_id_values,
metrics,
)
.await?;
}
@@ -400,10 +405,12 @@ pub async fn fetch_analytics(
let projects =
fetch_response_projects(&mut time_slices, &user, &pool, &redis).await?;
let users = fetch_response_users(&time_slices, &pool, &redis).await?;
Ok(web::Json(GetResponse {
metrics: time_slices,
projects,
users,
project_events,
}))
}
@@ -523,6 +530,43 @@ async fn fetch_response_projects(
.collect())
}
async fn fetch_response_users(
time_slices: &[TimeSlice],
pool: &PgPool,
redis: &RedisPool,
) -> Result<HashMap<UserId, User>, ApiError> {
let mut user_ids = HashSet::<database::models::DBUserId>::new();
for time_slice in time_slices {
for data in &time_slice.0 {
let AnalyticsData::Project(project) = data else {
continue;
};
if let ProjectMetrics::Revenue(revenue) = &project.metrics
&& let Some(user_id) = revenue.user_id
{
user_ids.insert(user_id.into());
}
}
}
let user_ids = user_ids.into_iter().collect::<Vec<_>>();
if user_ids.is_empty() {
return Ok(HashMap::new());
}
let users = DBUser::get_many_ids(&user_ids, pool, redis).await?;
Ok(users
.into_iter()
.map(|user| {
let user_id = UserId::from(user.id);
(user_id, User::from(user))
})
.collect())
}
fn filter_response_project_ids(
time_slices: &mut [TimeSlice],
visible_project_ids: &HashSet<DBProjectId>,
@@ -847,6 +891,7 @@ async fn filter_allowed_project_ids(
#[cfg(test)]
mod tests {
use crate::models::v3::users::{Badges, Role, UserCampaigns};
use rust_decimal::Decimal;
use serde_json::json;
@@ -965,11 +1010,13 @@ mod tests {
TimeSlice(vec![AnalyticsData::Project(ProjectAnalytics {
source_project: test_project_3,
metrics: ProjectMetrics::Revenue(ProjectRevenue {
user_id: None,
revenue: Decimal::new(20000, 2),
}),
})]),
],
projects: HashMap::new(),
users: HashMap::new(),
project_events: vec![],
};
let target = json!({
@@ -1002,4 +1049,90 @@ mod tests {
assert_eq!(serde_json::to_value(src).unwrap(), target);
}
#[test]
fn response_format_with_revenue_user() {
let test_project = ProjectId(123);
let test_user = UserId(456);
let created = DateTime::parse_from_rfc3339("2026-01-01T00:00:00Z")
.unwrap()
.with_timezone(&Utc);
let user = User {
id: test_user,
username: "revenue-user".into(),
avatar_url: None,
bio: None,
created,
role: Role::Developer,
badges: Badges::empty(),
campaigns: UserCampaigns { pride_26: None },
auth_providers: None,
email: None,
email_verified: None,
has_password: None,
has_totp: None,
payout_data: None,
stripe_customer_id: None,
allow_friend_requests: None,
moderation_notes: None,
github_id: None,
};
let src = GetResponse {
metrics: vec![TimeSlice(vec![AnalyticsData::Project(
ProjectAnalytics {
source_project: test_project,
metrics: ProjectMetrics::Revenue(ProjectRevenue {
user_id: Some(test_user),
revenue: Decimal::new(1234, 2),
}),
},
)])],
projects: HashMap::new(),
users: HashMap::from([(test_user, user)]),
project_events: vec![],
};
let mut target_users = serde_json::Map::new();
target_users.insert(
test_user.to_string(),
json!({
"id": test_user.to_string(),
"username": "revenue-user",
"avatar_url": null,
"bio": null,
"created": "2026-01-01T00:00:00Z",
"role": "developer",
"badges": 0,
"campaigns": {
"pride_26": null
},
"auth_providers": null,
"email": null,
"email_verified": null,
"has_password": null,
"has_totp": null,
"payout_data": null,
"stripe_customer_id": null,
"allow_friend_requests": null,
"github_id": null
}),
);
let target = json!({
"metrics": [
[
{
"source_project": test_project.to_string(),
"metric_kind": "revenue",
"user_id": test_user.to_string(),
"revenue": "12.34",
}
]
],
"projects": {},
"users": target_users,
"project_events": []
});
assert_eq!(serde_json::to_value(src).unwrap(), target);
}
}