use chrono::{Datelike, Duration, NaiveDate};
use serde_json::Value as Json;
use uuid::Uuid;
use crate::domain::event::{record_channel, stage_bus_event};
use crate::infrastructure::persistence::activity_repository::{
ActivityRepository, ActivityRow, NewActivityRow,
};
use crate::infrastructure::persistence::message_pipeline_repository::MessagePipelineRepository;
#[derive(Debug, thiserror::Error)]
pub enum ActivityError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("invalid input: {0}")]
Invalid(String),
#[error("not found: {0}")]
NotFound(String),
}
#[derive(Debug, Clone)]
pub struct ScheduleActivity {
pub res_model: String,
pub res_id: Uuid,
pub activity_type_id: Option<Uuid>,
pub summary: Option<String>,
pub note: Option<String>,
pub date_deadline: NaiveDate,
pub user_id: Uuid,
pub requested_user_id: Option<Uuid>,
pub chained_next_activity: Option<Json>,
}
pub struct ActivityWriteService {
pool: sqlx::PgPool,
}
impl ActivityWriteService {
pub fn new(pool: sqlx::PgPool) -> Self {
Self { pool }
}
pub async fn schedule(
&self,
cmd: ScheduleActivity,
today: NaiveDate,
) -> Result<Uuid, ActivityError> {
self.validate(&cmd).await?;
let mut tx = self.pool.begin().await?;
let id = self.insert(&mut tx, &cmd, today).await?;
stage_bus_event(
&mut tx,
"ActivityScheduled",
"MailActivity",
id,
record_channel(&cmd.res_model, cmd.res_id),
"ActivityScheduled",
serde_json::json!({
"activity_id": id, "res_model": cmd.res_model, "res_id": cmd.res_id,
"user_id": cmd.user_id, "date_deadline": cmd.date_deadline.to_string(),
}),
)
.await?;
tx.commit().await?;
Ok(id)
}
pub async fn schedule_plan(
&self,
plan_id: Uuid,
res_model: &str,
res_id: Uuid,
fallback_user_id: Uuid,
today: NaiveDate,
) -> Result<Vec<Uuid>, ActivityError> {
let mut tx = self.pool.begin().await?;
let templates = ActivityRepository::list_plan_templates(&mut tx, plan_id).await?;
if templates.is_empty() {
return Err(ActivityError::NotFound(format!("plan {plan_id} has no templates")));
}
let mut ids = Vec::new();
for t in &templates {
let user_id = t.user_id.unwrap_or(fallback_user_id);
let deadline = shift_date(today, t.delay_count.unwrap_or(0), t.delay_unit.as_deref());
let cmd = ScheduleActivity {
res_model: res_model.into(),
res_id,
activity_type_id: Some(t.activity_type_id),
summary: t.summary.clone(),
note: t.note.clone(),
date_deadline: deadline,
user_id,
requested_user_id: None,
chained_next_activity: None,
};
let id = self.insert(&mut tx, &cmd, deadline).await?;
stage_bus_event(
&mut tx,
"ActivityScheduled",
"MailActivity",
id,
record_channel(res_model, res_id),
"ActivityScheduled",
serde_json::json!({
"activity_id": id, "plan_id": plan_id,
"res_model": res_model, "res_id": res_id, "user_id": user_id,
"date_deadline": deadline.to_string(),
}),
)
.await?;
ids.push(id);
}
tx.commit().await?;
Ok(ids)
}
pub async fn action_done(
&self,
activity_id: Uuid,
feedback: Option<&str>,
today: NaiveDate,
) -> Result<(bool, Option<Uuid>), ActivityError> {
let mut tx = self.pool.begin().await?;
let activity = ActivityRepository::find_activity(
&mut tx, activity_id)
.await?
.ok_or_else(|| ActivityError::NotFound(format!("activity {activity_id}")))?;
let note_body = feedback.unwrap_or("Activity done");
let message_id = Uuid::new_v4();
MessagePipelineRepository::insert_mail_message(
&mut tx,
&crate::infrastructure::persistence::message_pipeline_repository::NewMailMessageRow {
id: message_id,
subject: None,
body: note_body,
message_type: "notification",
subtype_id: None,
is_internal: true,
author_id: Some(activity.user_id),
author_guest_id: None,
email_from: None,
message_id: None,
reply_to: None,
model: Some(&activity.res_model),
res_id: Some(activity.res_id),
record_name: None,
},
)
.await?;
let channel_key = record_channel(&activity.res_model, activity.res_id);
stage_bus_event(
&mut tx,
"MessagePosted",
"MailMessage",
message_id,
channel_key.clone(),
"MessagePosted",
serde_json::json!({
"message_id": message_id, "message_type": "notification",
"model": activity.res_model, "res_id": activity.res_id,
"activity_done": activity_id,
}),
)
.await?;
let newly_done = ActivityRepository::archive_done(&mut tx, activity_id).await?;
let mut next_id = None;
if newly_done {
if let Some(type_id) = activity.activity_type_id {
if let Some(atype) = ActivityRepository::find_activity_type(&mut tx, type_id).await? {
if atype.chaining_type == "trigger" {
next_id = self
.create_chained_next(&mut tx, &activity, &channel_key, today)
.await?;
}
}
}
stage_bus_event(
&mut tx,
"ActivityDone",
"MailActivity",
activity_id,
channel_key,
"ActivityDone",
serde_json::json!({
"activity_id": activity_id, "chained_next_activity_id": next_id,
}),
)
.await?;
}
tx.commit().await?;
Ok((newly_done, next_id))
}
async fn create_chained_next(
&self,
tx: &mut sqlx::PgConnection,
activity: &ActivityRow,
channel_key: &str,
today: NaiveDate,
) -> Result<Option<Uuid>, ActivityError> {
let Some(spec) = &activity.chained_next_activity else {
return Ok(None);
};
let deadline = spec["date_deadline"]
.as_str()
.and_then(|d| NaiveDate::parse_from_str(d, "%Y-%m-%d").ok())
.unwrap_or(today);
let cmd = ScheduleActivity {
res_model: activity.res_model.clone(),
res_id: activity.res_id,
activity_type_id: spec["activity_type_id"].as_str().and_then(|s| Uuid::parse_str(s).ok()),
summary: spec["summary"].as_str().map(|s| s.to_string()),
note: spec["note"].as_str().map(|s| s.to_string()),
date_deadline: deadline,
user_id: spec["user_id"]
.as_str()
.and_then(|s| Uuid::parse_str(s).ok())
.unwrap_or(activity.user_id),
requested_user_id: None,
chained_next_activity: None,
};
let id = self.insert(tx, &cmd, today).await?;
stage_bus_event(
tx,
"ActivityScheduled",
"MailActivity",
id,
channel_key.to_string(),
"ActivityScheduled",
serde_json::json!({
"activity_id": id, "chained_from": activity.id,
"res_model": cmd.res_model, "res_id": cmd.res_id,
"user_id": cmd.user_id, "date_deadline": cmd.date_deadline.to_string(),
}),
)
.await?;
Ok(Some(id))
}
async fn validate(&self, cmd: &ScheduleActivity) -> Result<(), ActivityError> {
if cmd.res_model.trim().is_empty() {
return Err(ActivityError::Invalid("res_model is required".into()));
}
Ok(())
}
async fn insert(
&self,
tx: &mut sqlx::PgConnection,
cmd: &ScheduleActivity,
_today: NaiveDate,
) -> Result<Uuid, ActivityError> {
let id = Uuid::new_v4();
let state = if cmd.date_deadline < _today {
"overdue"
} else if cmd.date_deadline == _today {
"today"
} else {
"planned"
};
ActivityRepository::insert_activity(
tx,
&NewActivityRow {
id,
res_model: &cmd.res_model,
res_id: cmd.res_id,
activity_type_id: cmd.activity_type_id,
summary: cmd.summary.as_deref(),
note: cmd.note.as_deref(),
date_deadline: cmd.date_deadline,
user_id: cmd.user_id,
requested_user_id: cmd.requested_user_id,
state,
chained_next_activity: cmd.chained_next_activity.as_ref(),
},
)
.await?;
Ok(id)
}
}
fn shift_date(base: NaiveDate, count: i32, unit: Option<&str>) -> NaiveDate {
match unit {
Some("weeks") => base + Duration::weeks(count as i64),
Some("months") => shift_months(base, count),
_ => base + Duration::days(count as i64),
}
}
fn shift_months(base: NaiveDate, count: i32) -> NaiveDate {
let months = base.year() * 12 + base.month() as i32 - 1 + count;
let year = months.div_euclid(12);
let month = (months.rem_euclid(12) + 1) as u32;
let day = base.day().min(days_in_month(year, month));
NaiveDate::from_ymd_opt(year, month, day).unwrap_or(base)
}
fn days_in_month(year: i32, month: u32) -> u32 {
match month {
1 | 3 | 5 | 7 | 8 | 10 | 12 => 31,
4 | 6 | 9 | 11 => 30,
2 => {
if (year % 4 == 0 && year % 100 != 0) || year % 400 == 0 {
29
} else {
28
}
}
_ => 31,
}
}