textsender_api PR / Rustfmt (pull_request) Successful in 49s
textsender_api PR / Check (pull_request) Successful in 1m22s
textsender_api PR / Clippy (pull_request) Successful in 2m5s
Rust Build / Test Suite (pull_request) Successful in 35s
Rust Build / Rustfmt (pull_request) Successful in 41s
Reviewed-on: phoenix/textsender_api#21
281 lines
8.1 KiB
Rust
281 lines
8.1 KiB
Rust
use sqlx::Row;
|
|
|
|
pub async fn insert(
|
|
pool: &sqlx::PgPool,
|
|
scheduled_message: &textsender_models::message::scheduling::ScheduledMessage,
|
|
user_id: &uuid::Uuid,
|
|
) -> Result<(uuid::Uuid, time::OffsetDateTime), sqlx::Error> {
|
|
match sqlx::query(
|
|
r#"
|
|
INSERT INTO "scheduled_messages" (scheduled, status, user_id)
|
|
VALUES($1, $2, $3) RETURNING id, created;
|
|
"#,
|
|
)
|
|
.bind(scheduled_message.scheduled)
|
|
.bind(&scheduled_message.status)
|
|
.bind(user_id)
|
|
.fetch_one(pool)
|
|
.await
|
|
{
|
|
Ok(row) => {
|
|
let id: uuid::Uuid = row.try_get("id")?;
|
|
let created: time::OffsetDateTime = row.try_get("created")?;
|
|
Ok((id, created))
|
|
}
|
|
Err(_) => Err(sqlx::Error::RowNotFound),
|
|
}
|
|
}
|
|
|
|
pub async fn get(
|
|
pool: &sqlx::PgPool,
|
|
id: &uuid::Uuid,
|
|
) -> Result<textsender_models::message::scheduling::ScheduledMessage, sqlx::Error> {
|
|
match sqlx::query(
|
|
r#"
|
|
SELECT id, scheduled, created, status, user_id FROM "scheduled_messages"
|
|
WHERE
|
|
id = $1;
|
|
"#,
|
|
)
|
|
.bind(id)
|
|
.fetch_one(pool)
|
|
.await
|
|
{
|
|
Ok(row) => match parse_row(&row).await {
|
|
Ok(scheduled_message) => Ok(scheduled_message),
|
|
Err(err) => Err(err),
|
|
},
|
|
Err(_) => Err(sqlx::Error::RowNotFound),
|
|
}
|
|
}
|
|
|
|
pub async fn get_with_user_id(
|
|
pool: &sqlx::PgPool,
|
|
user_id: &uuid::Uuid,
|
|
) -> Result<Vec<textsender_models::message::scheduling::ScheduledMessage>, sqlx::Error> {
|
|
match sqlx::query(
|
|
r#"
|
|
SELECT id, scheduled, created, status, user_id FROM "scheduled_messages"
|
|
WHERE
|
|
user_id = $1;
|
|
"#,
|
|
)
|
|
.bind(user_id)
|
|
.fetch_all(pool)
|
|
.await
|
|
{
|
|
Ok(rows) => {
|
|
let mut messages: Vec<textsender_models::message::scheduling::ScheduledMessage> =
|
|
Vec::new();
|
|
for row in rows {
|
|
match parse_row(&row).await {
|
|
Ok(scheduled_message) => {
|
|
messages.push(scheduled_message);
|
|
}
|
|
Err(err) => {
|
|
return Err(err);
|
|
}
|
|
}
|
|
}
|
|
|
|
Ok(messages)
|
|
}
|
|
Err(_) => Err(sqlx::Error::RowNotFound),
|
|
}
|
|
}
|
|
|
|
pub async fn fetch(
|
|
pool: &sqlx::PgPool,
|
|
) -> Result<textsender_models::message::scheduling::ScheduledMessage, sqlx::Error> {
|
|
match sqlx::query(
|
|
r#"
|
|
UPDATE "scheduled_messages"
|
|
SET status = $1
|
|
WHERE id = (
|
|
SELECT id FROM "scheduled_messages"
|
|
WHERE status = $2
|
|
ORDER BY id
|
|
FOR UPDATE SKIP LOCKED
|
|
LIMIT 1
|
|
)
|
|
RETURNING id, scheduled, created, status, user_id
|
|
"#,
|
|
)
|
|
.bind(textsender_models::message::scheduling::PROCESSING)
|
|
.bind(textsender_models::message::scheduling::READY)
|
|
.fetch_one(pool)
|
|
.await
|
|
{
|
|
Ok(row) => match parse_row(&row).await {
|
|
Ok(scheduled_message) => Ok(scheduled_message),
|
|
Err(err) => Err(err),
|
|
},
|
|
Err(err) => Err(err),
|
|
}
|
|
}
|
|
|
|
pub async fn update_status(
|
|
pool: &sqlx::PgPool,
|
|
id: &uuid::Uuid,
|
|
status: &str,
|
|
) -> Result<(), sqlx::Error> {
|
|
match sqlx::query(
|
|
r#"
|
|
UPDATE "scheduled_messages" SET status = $1 WHERE id = $2;
|
|
"#,
|
|
)
|
|
.bind(status)
|
|
.bind(id)
|
|
.execute(pool)
|
|
.await
|
|
{
|
|
Ok(_row) => Ok(()),
|
|
Err(err) => Err(err),
|
|
}
|
|
}
|
|
|
|
async fn parse_row(
|
|
row: &sqlx::postgres::PgRow,
|
|
) -> Result<textsender_models::message::scheduling::ScheduledMessage, sqlx::Error> {
|
|
let id: uuid::Uuid = row.try_get("id")?;
|
|
let scheduled: time::OffsetDateTime = row.try_get("scheduled")?;
|
|
let status: String = row.try_get("status")?;
|
|
let created: time::OffsetDateTime = row.try_get("created")?;
|
|
let user_id: uuid::Uuid = row.try_get("user_id")?;
|
|
|
|
Ok(textsender_models::message::scheduling::ScheduledMessage {
|
|
id,
|
|
scheduled: Some(scheduled),
|
|
created: Some(created),
|
|
status,
|
|
user_id,
|
|
})
|
|
}
|
|
|
|
pub mod sched_msg_event {
|
|
use sqlx::Row;
|
|
|
|
pub async fn insert(
|
|
pool: &sqlx::PgPool,
|
|
scheduled_message_event: &textsender_models::message::scheduling::ScheduledMessageEvent,
|
|
) -> Result<(uuid::Uuid, time::OffsetDateTime), sqlx::Error> {
|
|
match sqlx::query(
|
|
r#"
|
|
INSERT INTO "scheduled_message_events" (contact_id, message_id, scheduled_message_id)
|
|
VALUES($1, $2, $3) RETURNING id, created;
|
|
"#,
|
|
)
|
|
.bind(scheduled_message_event.contact_id)
|
|
.bind(scheduled_message_event.message_id)
|
|
.bind(scheduled_message_event.scheduled_message_id)
|
|
.fetch_one(pool)
|
|
.await
|
|
{
|
|
Ok(row) => {
|
|
let id: uuid::Uuid = row.try_get("id")?;
|
|
let created: time::OffsetDateTime = row.try_get("created")?;
|
|
Ok((id, created))
|
|
}
|
|
Err(_) => Err(sqlx::Error::RowNotFound),
|
|
}
|
|
}
|
|
|
|
pub async fn get(
|
|
pool: &sqlx::PgPool,
|
|
id: &uuid::Uuid,
|
|
) -> Result<textsender_models::message::scheduling::ScheduledMessageEvent, sqlx::Error> {
|
|
match sqlx::query(
|
|
r#"
|
|
SELECT id, contact_id, message_id, scheduled_message_id, created FROM "scheduled_message_events"
|
|
WHERE
|
|
id = $1;
|
|
"#,
|
|
)
|
|
.bind(id)
|
|
.fetch_one(pool)
|
|
.await
|
|
{
|
|
Ok(row) => match parse_row(&row).await {
|
|
Ok(scheduled_message_event) => Ok(scheduled_message_event),
|
|
Err(err) => Err(err),
|
|
},
|
|
Err(_) => Err(sqlx::Error::RowNotFound),
|
|
}
|
|
}
|
|
|
|
pub async fn get_with_scheduled_message_id(
|
|
pool: &sqlx::PgPool,
|
|
scheduled_message_id: &uuid::Uuid,
|
|
) -> Result<Vec<textsender_models::message::scheduling::ScheduledMessageEvent>, sqlx::Error>
|
|
{
|
|
match sqlx::query(
|
|
r#"
|
|
SELECT id, contact_id, message_id, scheduled_message_id, created FROM "scheduled_message_events"
|
|
WHERE
|
|
scheduled_message_id = $1;
|
|
"#,
|
|
)
|
|
.bind(scheduled_message_id)
|
|
.fetch_all(pool)
|
|
.await
|
|
{
|
|
Ok(rows) => {
|
|
let mut scheduled_message_events: Vec<textsender_models::message::scheduling::ScheduledMessageEvent> =
|
|
Vec::new();
|
|
for row in rows {
|
|
match parse_row(&row).await {
|
|
Ok(scheduled_message_event) => {
|
|
scheduled_message_events.push(scheduled_message_event);
|
|
}
|
|
Err(err) => {
|
|
return Err(err);
|
|
}
|
|
}
|
|
}
|
|
|
|
Ok(scheduled_message_events)
|
|
}
|
|
Err(_) => Err(sqlx::Error::RowNotFound),
|
|
}
|
|
}
|
|
|
|
pub async fn delete(
|
|
pool: &sqlx::PgPool,
|
|
scheduled_message_event: &textsender_models::message::scheduling::ScheduledMessageEvent,
|
|
) -> Result<(), sqlx::Error> {
|
|
match sqlx::query(
|
|
r#"
|
|
DELETE FROM "scheduled_message_events"
|
|
WHERE id = $1
|
|
"#,
|
|
)
|
|
.bind(scheduled_message_event.id)
|
|
.execute(pool)
|
|
.await
|
|
{
|
|
Ok(_) => Ok(()),
|
|
Err(err) => Err(err),
|
|
}
|
|
}
|
|
|
|
async fn parse_row(
|
|
row: &sqlx::postgres::PgRow,
|
|
) -> Result<textsender_models::message::scheduling::ScheduledMessageEvent, sqlx::Error> {
|
|
let id: uuid::Uuid = row.try_get("id")?;
|
|
let contact_id: uuid::Uuid = row.try_get("contact_id")?;
|
|
let message_id: uuid::Uuid = row.try_get("message_id")?;
|
|
let scheduled_message_id: uuid::Uuid = row.try_get("scheduled_message_id")?;
|
|
let created: time::OffsetDateTime = row.try_get("created")?;
|
|
|
|
Ok(
|
|
textsender_models::message::scheduling::ScheduledMessageEvent {
|
|
id,
|
|
contact_id,
|
|
message_id,
|
|
scheduled_message_id,
|
|
created: Some(created),
|
|
},
|
|
)
|
|
}
|
|
}
|