A lot of changes #12
+1
-1
@@ -8,6 +8,6 @@ DB_MAIN_SSLMODE=disable
|
||||
DATABASE_URL=postgres://${DB_MAIN_USER}:${DB_MAIN_PASSWORD}@${DB_MAIN_HOST}:${DB_MAIN_PORT}/${DB_MAIN_NAME}
|
||||
TWILIO_AUTH_SID=9M438C93R943U4329MCU43C34U
|
||||
TWILIO_SERVICE_SID=9M4J3X8439U398NUVT3342MC349C348T
|
||||
TWILIO_AUTH_TOKEN="f4a1f2b0b79ea3735078c2d8ee9684e1"
|
||||
TWILIO_AUTH_TOKEN="MA08HJp9tYtkfs72m8JIIROofRVm74oR0ixEYBHGYpB9"
|
||||
TWILIO_PHONE_NUMBER=+10123456789
|
||||
ALLOWED_ORIGINS="http://textsender.com"
|
||||
|
||||
+1
-1
@@ -8,6 +8,6 @@ DB_MAIN_SSLMODE=disable
|
||||
DATABASE_URL=postgres://${DB_MAIN_USER}:${DB_MAIN_PASSWORD}@${DB_MAIN_HOST}:${DB_MAIN_PORT}/${DB_MAIN_NAME}
|
||||
TWILIO_AUTH_SID=9M438C93R943U4329MCU43C34U
|
||||
TWILIO_SERVICE_SID=9M4J3X8439U398NUVT3342MC349C348T
|
||||
TWILIO_AUTH_TOKEN="f4a1f2b0b79ea3735078c2d8ee9684e1"
|
||||
TWILIO_AUTH_TOKEN="MA08HJp9tYtkfs72m8JIIROofRVm74oR0ixEYBHGYpB9"
|
||||
TWILIO_PHONE_NUMBER=+10123456789
|
||||
ALLOWED_ORIGINS="http://textsender.com"
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
name: Rust Build
|
||||
|
||||
on:
|
||||
pull_request:
|
||||
branches:
|
||||
- next
|
||||
|
||||
|
||||
concurrency:
|
||||
group: ${{ github.workflow }}-${{ github.ref }}
|
||||
cancel-in-progress: true
|
||||
|
||||
jobs:
|
||||
test:
|
||||
name: Test Suite
|
||||
runs-on: ubuntu-24.04
|
||||
# --- Add database service definition ---
|
||||
services:
|
||||
postgres:
|
||||
image: postgres:18.4-alpine
|
||||
env:
|
||||
# Use secrets for DB init, with fallbacks for flexibility
|
||||
POSTGRES_USER: ${{ secrets.DB_TEST_USER || 'testuser' }}
|
||||
POSTGRES_PASSWORD: ${{ secrets.DB_TEST_PASSWORD || 'testpassword' }}
|
||||
POSTGRES_DB: ${{ secrets.DB_TEST_NAME || 'testdb' }}
|
||||
POSTGRES_PORT: ${{ secrets.DB_PORT || 5432 }}
|
||||
# Options to wait until the database is ready
|
||||
options: >-
|
||||
--health-cmd pg_isready
|
||||
--health-interval 10s
|
||||
--health-timeout 5s
|
||||
--health-retries 5
|
||||
|
||||
steps:
|
||||
- uses: actions/checkout@v5
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
toolchain: 1.96
|
||||
# --- Add this step for explicit verification ---
|
||||
- name: Verify Docker Environment
|
||||
run: |
|
||||
echo "Runner User Info:"
|
||||
id
|
||||
echo "Checking Docker Version:"
|
||||
docker --version
|
||||
echo "Checking Docker Daemon Status (info):"
|
||||
docker info
|
||||
echo "Checking Docker Daemon Status (ps):"
|
||||
docker ps -a
|
||||
echo "Docker environment check complete."
|
||||
# NOTE: Do NOT use continue-on-error here.
|
||||
# If Docker isn't working as expected, the job SHOULD fail here.
|
||||
- name: Run tests
|
||||
env:
|
||||
# Define DATABASE_URL for tests to use
|
||||
DATABASE_URL: postgresql://${{ secrets.DB_TEST_USER || 'testuser' }}:${{ secrets.DB_TEST_PASSWORD || 'testpassword' }}@postgres:${{ secrets.DB_PORT || 5432 }}/${{ secrets.DB_TEST_NAME || 'testdb' }}
|
||||
RUST_LOG: info # Optional: configure test log level
|
||||
SECRET_MAIN_KEY: ${{ secrets.SECRET_KEY }}
|
||||
# Make SSH agent available if tests fetch private dependencies
|
||||
SSH_AUTH_SOCK: ${{ env.SSH_AUTH_SOCK }}
|
||||
ENABLE_REGISTRATION: 'TRUE'
|
||||
run: |
|
||||
mkdir -p ~/.ssh
|
||||
echo "${{ secrets.MYREPO_TOKEN }}" > ~/.ssh/textsender-models_deploy_key
|
||||
chmod 600 ~/.ssh/textsender-models_deploy_key
|
||||
ssh-keyscan ${{ secrets.MY_HOST }} >> ~/.ssh/known_hosts
|
||||
|
||||
eval $(ssh-agent -s)
|
||||
ssh-add -v ~/.ssh/textsender-models_deploy_key
|
||||
|
||||
cargo test
|
||||
|
||||
fmt:
|
||||
name: Rustfmt
|
||||
runs-on: ubuntu-24.04
|
||||
steps:
|
||||
- uses: actions/checkout@v5
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
toolchain: 1.96
|
||||
- run: rustup component add rustfmt
|
||||
- uses: Swatinem/rust-cache@v2
|
||||
- run: |
|
||||
mkdir -p ~/.ssh
|
||||
echo "${{ secrets.MYREPO_TOKEN }}" > ~/.ssh/textsender-models_deploy_key
|
||||
chmod 600 ~/.ssh/textsender-models_deploy_key
|
||||
ssh-keyscan ${{ secrets.MY_HOST }} >> ~/.ssh/known_hosts
|
||||
|
||||
eval $(ssh-agent -s)
|
||||
ssh-add -v ~/.ssh/textsender-models_deploy_key
|
||||
cargo fmt --all -- --check
|
||||
@@ -37,7 +37,7 @@ jobs:
|
||||
# --- Add database service definition ---
|
||||
services:
|
||||
postgres:
|
||||
image: postgres:18.3-alpine
|
||||
image: postgres:18.4-alpine
|
||||
env:
|
||||
# Use secrets for DB init, with fallbacks for flexibility
|
||||
POSTGRES_USER: ${{ secrets.DB_TEST_USER || 'testuser' }}
|
||||
|
||||
Generated
+1248
-297
File diff suppressed because it is too large
Load Diff
+4
-7
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "textsender_api"
|
||||
version = "0.1.7"
|
||||
version = "0.2.0"
|
||||
edition = "2024"
|
||||
rust-version = "1.96"
|
||||
|
||||
@@ -15,19 +15,16 @@ tower = { version = "0.5.3", features = ["full"] }
|
||||
tower-http = { version = "0.6.11", features = ["cors", "timeout"] }
|
||||
tracing-subscriber = "0.3.23"
|
||||
futures = { version = "0.3.32" }
|
||||
mime_guess = { version = "2.0.5" }
|
||||
uuid = { version = "1.23.3", features = ["v4", "serde"] }
|
||||
sqlx = { version = "0.8.6", features = ["postgres", "runtime-tokio-native-tls", "time", "uuid"] }
|
||||
sqlx = { version = "0.9.0", features = ["runtime-tokio", "tls-native-tls", "postgres", "time", "uuid"] }
|
||||
time = { version = "0.3.49", features = ["formatting", "macros", "parsing", "serde"] }
|
||||
thiserror = "2.0.18"
|
||||
base64 = "0.22.1"
|
||||
jsonwebtoken = { version = "10.4.0", features = ["rust_crypto"] }
|
||||
josekit = { version = "0.10.3" }
|
||||
utoipa = { version = "5.5.0", features = ["axum_extras"] }
|
||||
utoipa-swagger-ui = { version = "9.0.2", features = ["axum"] }
|
||||
textsender_models = { git = "ssh://git@git.kundeng.us/phoenix/textsender_models.git", tag = "v0.4.0" }
|
||||
textsender_models = { git = "ssh://git@git.kundeng.us/phoenix/textsender_models.git", tag = "v0.4.10" }
|
||||
swoosh = { git = "ssh://git@git.kundeng.us/phoenix/swoosh.git", tag = "v0.4.2.0" }
|
||||
|
||||
[dev-dependencies]
|
||||
common-multipart-rfc7578 = { version = "0.7.0" }
|
||||
url = { version = "2.5.8" }
|
||||
tempfile = { version = "3.27.0" }
|
||||
|
||||
+17
-17
@@ -45,28 +45,28 @@ services:
|
||||
# - ./path/to/your/web-api-repo:/app
|
||||
|
||||
# --- catapult service ---
|
||||
# catapult:
|
||||
# build:
|
||||
# context: ../catapult
|
||||
# ssh: ["default"]
|
||||
# dockerfile: Dockerfile
|
||||
# container_name: catapult
|
||||
# restart: unless-stopped
|
||||
# env_file:
|
||||
# - ../catapult/.env
|
||||
# depends_on:
|
||||
# - api
|
||||
# - main_db
|
||||
# - auth_api
|
||||
# - auth_db
|
||||
# networks:
|
||||
# - textsender_api-network
|
||||
catapult:
|
||||
build:
|
||||
context: ../catapult
|
||||
ssh: ["default"]
|
||||
dockerfile: Dockerfile
|
||||
container_name: catapult
|
||||
restart: unless-stopped
|
||||
env_file:
|
||||
- ../catapult/.env
|
||||
depends_on:
|
||||
- api
|
||||
- main_db
|
||||
- auth_api
|
||||
- auth_db
|
||||
networks:
|
||||
- textsender_api-network
|
||||
|
||||
|
||||
# PostgreSQL Database Service
|
||||
# --- textsender_api web api db ---
|
||||
main_db:
|
||||
image: postgres:18.3-alpine # Use an official Postgres image (Alpine variant is smaller)
|
||||
image: postgres:18.4-alpine # Use an official Postgres image (Alpine variant is smaller)
|
||||
container_name: textsender_api_db # Optional: Give the container a specific name
|
||||
environment:
|
||||
# These MUST match the user, password, and database name in the DATABASE_URL above
|
||||
|
||||
@@ -10,3 +10,36 @@ CREATE TABLE IF NOT EXISTS "contacts" (
|
||||
lastname TEXT NULL,
|
||||
nickname TEXT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS "messages" (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
content TEXT NOT NULL,
|
||||
user_id UUID NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS "scheduled_messages" (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
scheduled timestamptz NOT NULL,
|
||||
created timestamptz DEFAULT now(),
|
||||
status TEXT CHECK (status IN ('PENDING', 'READY', 'PROCESSING', 'DONE')),
|
||||
user_id UUID NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS "scheduled_message_events" (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
contact_id UUID NOT NULL,
|
||||
message_id UUID NOT NULL,
|
||||
scheduled_message_id UUID NOT NULL,
|
||||
created timestamptz DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS "message_event_responses" (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
scheduled_message_event_id UUID NULL,
|
||||
response JSONB NOT NULL,
|
||||
user_id UUID NOT NULL,
|
||||
sent timestamptz NOT NULL,
|
||||
contact_id UUID NULL,
|
||||
message_id UUID NULL,
|
||||
status TEXT CHECK (status IN ('INSTANT', 'SCHEDULED'))
|
||||
);
|
||||
|
||||
@@ -1,49 +0,0 @@
|
||||
-- CREATE EXTENSION IF NOT EXISTS "uuid-ossp";
|
||||
--
|
||||
-- DROP TABLE IF EXISTS contacts CASCADE;
|
||||
-- DROP TABLE IF EXISTS messages CASCADE;
|
||||
-- DROP TABLE IF EXISTS scheduled_messages CASCADE;
|
||||
-- DROP TABLE IF EXISTS scheduled_message_events CASCADE;
|
||||
-- DROP TABLE IF EXISTS message_event_responses CASCADE;
|
||||
--
|
||||
-- CREATE TABLE IF NOT EXISTS contacts (
|
||||
-- id UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
|
||||
-- phone_number TEXT NOT NULL,
|
||||
-- user_id UUID NOT NULL,
|
||||
-- first_name TEXT NULL,
|
||||
-- last_name TEXT NULL,
|
||||
-- nickname TEXT NULL
|
||||
-- );
|
||||
--
|
||||
-- CREATE TABLE IF NOT EXISTS messages (
|
||||
-- id UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
|
||||
-- content TEXT NOT NULL,
|
||||
-- user_id UUID NOT NULL
|
||||
-- );
|
||||
--
|
||||
-- CREATE TABLE IF NOT EXISTS scheduled_messages (
|
||||
-- id UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
|
||||
-- scheduled timestamptz NOT NULL,
|
||||
-- created timestamptz DEFAULT now(),
|
||||
-- status TEXT CHECK (status IN ('PENDING', 'READY', 'PROCESSING', 'DONE')),
|
||||
-- user_id UUID NOT NULL
|
||||
-- );
|
||||
--
|
||||
-- CREATE TABLE IF NOT EXISTS scheduled_message_events (
|
||||
-- id UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
|
||||
-- contact_id UUID NOT NULL,
|
||||
-- message_id UUID NOT NULL,
|
||||
-- scheduled_message_id UUID NOT NULL,
|
||||
-- created timestamptz DEFAULT now()
|
||||
-- );
|
||||
--
|
||||
-- CREATE TABLE IF NOT EXISTS message_event_responses (
|
||||
-- id UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
|
||||
-- scheduled_message_event_id UUID NULL,
|
||||
-- response JSONB NOT NULL,
|
||||
-- user_id UUID NOT NULL,
|
||||
-- sent timestamptz NOT NULL,
|
||||
-- contact_id UUID NULL,
|
||||
-- message_id UUID NULL,
|
||||
-- status TEXT CHECK (status IN ('INSTANT', 'SCHEDULED'))
|
||||
-- );
|
||||
@@ -171,7 +171,7 @@ pub mod endpoint {
|
||||
}
|
||||
}
|
||||
|
||||
// Endpoint to get songs
|
||||
/// Endpoint to get Contacts
|
||||
#[utoipa::path(
|
||||
get,
|
||||
path = super::super::endpoints::GET_CONTACT,
|
||||
|
||||
@@ -0,0 +1,171 @@
|
||||
pub mod request {
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::ToSchema;
|
||||
|
||||
#[derive(Debug, Deserialize, Serialize, ToSchema)]
|
||||
pub struct InstantMessageRequest {
|
||||
pub contact_ids: Vec<uuid::Uuid>,
|
||||
pub message_id: uuid::Uuid,
|
||||
pub user_id: uuid::Uuid,
|
||||
}
|
||||
|
||||
impl InstantMessageRequest {
|
||||
pub fn is_valid(&self) -> bool {
|
||||
!self.contact_ids.is_empty() || !self.message_id.is_nil() || !self.user_id.is_nil()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub mod response {
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::ToSchema;
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct InstantMessageResponse {
|
||||
pub message: String,
|
||||
pub data: Vec<textsender_models::message::event::MessageEventResponse>,
|
||||
}
|
||||
}
|
||||
|
||||
pub mod endpoint {
|
||||
use crate::repo::contact as contact_repo;
|
||||
use crate::repo::message as message_repo;
|
||||
use message_repo::event as event_repo;
|
||||
|
||||
/// Endpoint to create Message
|
||||
#[utoipa::path(
|
||||
post,
|
||||
path = super::super::endpoints::ADD_MESSAGE,
|
||||
request_body(
|
||||
content = super::request::InstantMessageRequest,
|
||||
description = "Data needed to create a Message",
|
||||
content_type = "application/json"
|
||||
),
|
||||
responses(
|
||||
(status = 201, description = "Message created", body = super::response::InstantMessageResponse),
|
||||
(status = 400, description = "Error", body = super::response::InstantMessageResponse),
|
||||
(status = 500, description = "Error creating Message", body = super::response::InstantMessageResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn send_message(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::Json(payload): axum::Json<super::request::InstantMessageRequest>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::InstantMessageResponse>,
|
||||
) {
|
||||
let mut response = super::response::InstantMessageResponse::default();
|
||||
|
||||
if payload.is_valid() {
|
||||
let mut contacts: Vec<textsender_models::contact::Contact> = Vec::new();
|
||||
|
||||
for contact_id in &payload.contact_ids {
|
||||
match contact_repo::get(&pool, contact_id).await {
|
||||
Ok(contact) => {
|
||||
contacts.push(contact);
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Invalid");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if !contacts.is_empty() {
|
||||
println!("Valid contacts");
|
||||
|
||||
match message_repo::get(&pool, &payload.message_id).await {
|
||||
Ok(message) => {
|
||||
println!("Valid message");
|
||||
let t_config =
|
||||
match textsender_models::config::auxiliary::load_config().await {
|
||||
Ok(config) => config,
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Error getting config");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
};
|
||||
let mut results: Vec<
|
||||
textsender_models::message::event::MessageEventResponse,
|
||||
> = Vec::new();
|
||||
let pp = swoosh::twilio::types::Parameters {
|
||||
schedule: false,
|
||||
schedule_at: None,
|
||||
};
|
||||
|
||||
for contact in &contacts {
|
||||
match swoosh::twilio::api::send_message(
|
||||
&message, contact, &pp, &t_config,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(result) => {
|
||||
let val = swoosh::twilio::api::response_to_json(result).await;
|
||||
let mer = textsender_models::message::event::MessageEventResponse {
|
||||
user_id: payload.user_id,
|
||||
// TODO: Make this more accurate
|
||||
sent: Some(time::OffsetDateTime::now_utc()),
|
||||
status: String::from(textsender_models::message::event::MESSAGE_EVENT_RESPONSE_STATUS_INSTANT),
|
||||
contact_id: contact.id.unwrap(),
|
||||
message_id: message.id.unwrap(),
|
||||
response: val,
|
||||
..Default::default()
|
||||
};
|
||||
results.push(mer);
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if results.is_empty() {
|
||||
eprintln!("Error sending");
|
||||
(
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
)
|
||||
} else {
|
||||
for result in &mut results {
|
||||
match event_repo::insert(&pool, result).await {
|
||||
Ok(i) => {
|
||||
result.id = i;
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
response.data = results;
|
||||
response.message = String::from(super::super::response::SUCCESSFUL);
|
||||
(axum::http::StatusCode::OK, axum::Json(response))
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
} else {
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
} else {
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,167 @@
|
||||
pub mod request {
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::ToSchema;
|
||||
|
||||
#[derive(Debug, Deserialize, Serialize, ToSchema)]
|
||||
pub struct RecordMessageEventRequest {
|
||||
pub scheduled_message_event_id: uuid::Uuid,
|
||||
pub response: serde_json::Value,
|
||||
pub user_id: uuid::Uuid,
|
||||
#[serde(with = "time::serde::rfc3339::option")]
|
||||
pub sent: Option<time::OffsetDateTime>,
|
||||
pub status: String,
|
||||
pub contact_id: uuid::Uuid,
|
||||
pub message_id: uuid::Uuid,
|
||||
}
|
||||
|
||||
impl RecordMessageEventRequest {
|
||||
pub fn is_valid(&self) -> bool {
|
||||
!self.scheduled_message_event_id.is_nil()
|
||||
|| !self.response.is_null()
|
||||
|| !self.user_id.is_nil()
|
||||
|| self.sent.is_some()
|
||||
|| !self.status.is_empty()
|
||||
|| !self.contact_id.is_nil()
|
||||
|| !self.message_id.is_nil()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, serde::Deserialize, serde::Serialize, utoipa::ToSchema)]
|
||||
pub struct GetMessageEventResponseParams {
|
||||
pub id: Option<uuid::Uuid>,
|
||||
pub user_id: Option<uuid::Uuid>,
|
||||
}
|
||||
}
|
||||
|
||||
pub mod response {
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::ToSchema;
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct RecordMessageEventResponse {
|
||||
pub message: String,
|
||||
pub data: Vec<textsender_models::message::event::MessageEventResponse>,
|
||||
}
|
||||
|
||||
pub use RecordMessageEventResponse as GetMERResponse;
|
||||
}
|
||||
|
||||
pub mod endpoint {
|
||||
use crate::repo::message::event as event_repo;
|
||||
|
||||
/// Endpoint to create a Message Event Response
|
||||
#[utoipa::path(
|
||||
post,
|
||||
path = crate::caller::endpoints::RECORD_MESSAGE_EVENT_RESPONSE,
|
||||
request_body(
|
||||
content = super::request::RecordMessageEventRequest,
|
||||
description = "Data needed to create a Message Event Response",
|
||||
content_type = "application/json"
|
||||
),
|
||||
responses(
|
||||
(status = 201, description = "Message Event Response created", body = super::response::RecordMessageEventResponse),
|
||||
(status = 400, description = "Error", body = super::response::RecordMessageEventResponse),
|
||||
(status = 500, description = "Error creating Message Event Response", body = super::response::RecordMessageEventResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn record_message_event_response(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::Json(payload): axum::Json<super::request::RecordMessageEventRequest>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::RecordMessageEventResponse>,
|
||||
) {
|
||||
let mut response = super::response::RecordMessageEventResponse::default();
|
||||
|
||||
if payload.is_valid() {
|
||||
let mut mer = textsender_models::message::event::MessageEventResponse {
|
||||
scheduled_message_event_id: Some(payload.scheduled_message_event_id),
|
||||
user_id: payload.user_id,
|
||||
sent: payload.sent,
|
||||
status: payload.status,
|
||||
contact_id: payload.contact_id,
|
||||
message_id: payload.message_id,
|
||||
response: payload.response,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
match event_repo::insert(&pool, &mer).await {
|
||||
Ok(id) => {
|
||||
mer.id = id;
|
||||
response.message = String::from(super::super::super::response::SUCCESSFUL);
|
||||
response.data.push(mer);
|
||||
(axum::http::StatusCode::CREATED, axum::Json(response))
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
} else {
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
|
||||
/// Endpoint to get Messages
|
||||
#[utoipa::path(
|
||||
get,
|
||||
path = crate::caller::endpoints::GET_MESSAGE_EVENT_RESPONSE,
|
||||
params(
|
||||
("id" = uuid::Uuid, Path, description = "Id of Message Event Response"),
|
||||
("user_id" = uuid::Uuid, Path, description = "User Id associated with the Message Event Response")
|
||||
),
|
||||
responses(
|
||||
(status = 200, description = "Message Event Response found", body = super::response::GetMERResponse),
|
||||
(status = 400, description = "Error getting Message Event Response", body = super::response::GetMERResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn get_record_message_event_response(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::extract::Query(params): axum::extract::Query<
|
||||
super::request::GetMessageEventResponseParams,
|
||||
>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::GetMERResponse>,
|
||||
) {
|
||||
let mut response = super::response::GetMERResponse::default();
|
||||
|
||||
let messages = match params.id {
|
||||
Some(id) => match event_repo::get(&pool, &id).await {
|
||||
Ok(message) => {
|
||||
vec![message]
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
},
|
||||
None => match params.user_id {
|
||||
Some(user_id) => match event_repo::get_with_user_id(&pool, &user_id).await {
|
||||
Ok(messages) => messages,
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
},
|
||||
None => {
|
||||
response.message = String::from("Invalid parameter");
|
||||
return (axum::http::StatusCode::BAD_REQUEST, axum::Json(response));
|
||||
}
|
||||
},
|
||||
};
|
||||
|
||||
response.data = messages;
|
||||
response.message = String::from(super::super::super::response::SUCCESSFUL);
|
||||
|
||||
(axum::http::StatusCode::OK, axum::Json(response))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,154 @@
|
||||
pub mod event;
|
||||
pub mod scheduling;
|
||||
|
||||
pub mod request {
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::ToSchema;
|
||||
|
||||
#[derive(Debug, Deserialize, Serialize, ToSchema)]
|
||||
pub struct AddMessageRequest {
|
||||
pub content: String,
|
||||
pub user_id: uuid::Uuid,
|
||||
}
|
||||
|
||||
impl AddMessageRequest {
|
||||
pub fn is_valid(&self) -> bool {
|
||||
!self.content.is_empty() || !self.user_id.is_nil()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, serde::Deserialize, serde::Serialize, utoipa::ToSchema)]
|
||||
pub struct GetMessageParams {
|
||||
pub id: Option<uuid::Uuid>,
|
||||
pub user_id: Option<uuid::Uuid>,
|
||||
}
|
||||
}
|
||||
|
||||
pub mod response {
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::ToSchema;
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct AddMessageResponse {
|
||||
pub message: String,
|
||||
pub data: Vec<textsender_models::message::Message>,
|
||||
}
|
||||
|
||||
pub use AddMessageResponse as GetMessageResponse;
|
||||
}
|
||||
|
||||
pub mod endpoint {
|
||||
use crate::repo::message as message_repo;
|
||||
|
||||
/// Endpoint to create Message
|
||||
#[utoipa::path(
|
||||
post,
|
||||
path = super::super::endpoints::ADD_MESSAGE,
|
||||
request_body(
|
||||
content = super::request::AddMessageRequest,
|
||||
description = "Data needed to create a Message",
|
||||
content_type = "application/json"
|
||||
),
|
||||
responses(
|
||||
(status = 201, description = "Message created", body = super::response::AddMessageResponse),
|
||||
(status = 400, description = "Error", body = super::response::AddMessageResponse),
|
||||
(status = 500, description = "Error creating Message", body = super::response::AddMessageResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn create_message(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::Json(payload): axum::Json<super::request::AddMessageRequest>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::AddMessageResponse>,
|
||||
) {
|
||||
let mut response = super::response::AddMessageResponse::default();
|
||||
|
||||
if payload.is_valid() {
|
||||
let mut msg = textsender_models::message::Message {
|
||||
content: payload.content,
|
||||
user_id: Some(payload.user_id),
|
||||
..Default::default()
|
||||
};
|
||||
match message_repo::insert(&pool, &msg, &payload.user_id).await {
|
||||
Ok(message_id) => {
|
||||
msg.id = Some(message_id);
|
||||
response.data.push(msg);
|
||||
response.message = String::from("Message created");
|
||||
(axum::http::StatusCode::CREATED, axum::Json(response))
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Contact not created");
|
||||
|
||||
(
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
|
||||
/// Endpoint to get Messages
|
||||
#[utoipa::path(
|
||||
get,
|
||||
path = super::super::endpoints::GET_MESSAGE,
|
||||
params(
|
||||
("id" = uuid::Uuid, Path, description = "Id of Message"),
|
||||
("user_id" = uuid::Uuid, Path, description = "User Id associated with the Message")
|
||||
),
|
||||
responses(
|
||||
(status = 200, description = "Songs found", body = super::response::GetMessageResponse),
|
||||
(status = 400, description = "Error getting songs", body = super::response::GetMessageResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn get_messages(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::extract::Query(params): axum::extract::Query<super::request::GetMessageParams>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::GetMessageResponse>,
|
||||
) {
|
||||
let mut response = super::response::GetMessageResponse::default();
|
||||
|
||||
let messages = match params.id {
|
||||
Some(id) => match message_repo::get(&pool, &id).await {
|
||||
Ok(message) => {
|
||||
vec![message]
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
},
|
||||
None => match params.user_id {
|
||||
Some(user_id) => match message_repo::get_with_user_id(&pool, &user_id).await {
|
||||
Ok(messages) => messages,
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
},
|
||||
None => {
|
||||
response.message = String::from("Invalid parameter");
|
||||
return (axum::http::StatusCode::BAD_REQUEST, axum::Json(response));
|
||||
}
|
||||
},
|
||||
};
|
||||
|
||||
response.data = messages;
|
||||
response.message = String::from(super::super::response::SUCCESSFUL);
|
||||
|
||||
(axum::http::StatusCode::OK, axum::Json(response))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,549 @@
|
||||
pub mod request {
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::ToSchema;
|
||||
|
||||
#[derive(Debug, Deserialize, Serialize, ToSchema)]
|
||||
pub struct ScheduleMessageRequest {
|
||||
#[serde(with = "time::serde::rfc3339::option")]
|
||||
pub scheduled: Option<time::OffsetDateTime>,
|
||||
pub status: String,
|
||||
pub user_id: uuid::Uuid,
|
||||
}
|
||||
|
||||
impl ScheduleMessageRequest {
|
||||
pub fn is_valid(&self) -> bool {
|
||||
self.scheduled.is_some() || !self.status.is_empty() || !self.user_id.is_nil()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, serde::Deserialize, serde::Serialize, utoipa::ToSchema)]
|
||||
pub struct GetScheduledMessageParams {
|
||||
pub id: Option<uuid::Uuid>,
|
||||
pub user_id: Option<uuid::Uuid>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, serde::Deserialize, serde::Serialize, utoipa::ToSchema)]
|
||||
pub struct CreateScheduleMessageEventRequest {
|
||||
pub contact_id: uuid::Uuid,
|
||||
pub message_id: uuid::Uuid,
|
||||
pub scheduled_message_id: uuid::Uuid,
|
||||
}
|
||||
|
||||
impl CreateScheduleMessageEventRequest {
|
||||
pub fn is_valid(&self) -> bool {
|
||||
!self.contact_id.is_nil()
|
||||
|| !self.message_id.is_nil()
|
||||
|| !self.scheduled_message_id.is_nil()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, serde::Deserialize, serde::Serialize, utoipa::ToSchema)]
|
||||
pub struct GetScheduledMessageEventParams {
|
||||
pub id: Option<uuid::Uuid>,
|
||||
pub scheduled_message_id: Option<uuid::Uuid>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct UpdateScheduledMessageStatusRequest {
|
||||
pub status: String,
|
||||
pub scheduled_message_id: uuid::Uuid,
|
||||
}
|
||||
|
||||
impl UpdateScheduledMessageStatusRequest {
|
||||
pub fn is_valid(&self) -> bool {
|
||||
if !self.status.is_empty() || !self.scheduled_message_id.is_nil() {
|
||||
self.status == textsender_models::message::scheduling::PENDING
|
||||
|| self.status == textsender_models::message::scheduling::READY
|
||||
|| self.status == textsender_models::message::scheduling::PROCESSING
|
||||
|| self.status == textsender_models::message::scheduling::DONE
|
||||
} else {
|
||||
false
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub mod response {
|
||||
use serde::{Deserialize, Serialize};
|
||||
use utoipa::ToSchema;
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct ScheduleMessageResponse {
|
||||
pub message: String,
|
||||
pub data: Vec<textsender_models::message::scheduling::ScheduledMessage>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct GetScheduledMessageResponse {
|
||||
pub message: String,
|
||||
pub data: Vec<textsender_models::message::scheduling::ScheduledMessage>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct CreateScheduleMessageEventResponse {
|
||||
pub message: String,
|
||||
pub data: Vec<textsender_models::message::scheduling::ScheduledMessageEvent>,
|
||||
}
|
||||
|
||||
pub use CreateScheduleMessageEventResponse as GetScheduledMessageEventResponse;
|
||||
|
||||
pub use CreateScheduleMessageEventResponse as DeleteScheduledMessageEventResponse;
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct UpdateScheduledMessageStatusResponse {
|
||||
pub message: String,
|
||||
pub data: Vec<ChangedStatus>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize, ToSchema)]
|
||||
pub struct ChangedStatus {
|
||||
pub old_status: String,
|
||||
}
|
||||
|
||||
pub use ScheduleMessageResponse as FetchScheduledMessageResponse;
|
||||
}
|
||||
|
||||
pub mod endpoint {
|
||||
use crate::repo::contact as contact_repo;
|
||||
use crate::repo::message as message_repo;
|
||||
use message_repo::scheduling as scheduling_repo;
|
||||
|
||||
/// Endpoint to create a Scheudled Message
|
||||
#[utoipa::path(
|
||||
post,
|
||||
path = crate::caller::endpoints::SCHEDULE_MESSAGE,
|
||||
request_body(
|
||||
content = super::request::ScheduleMessageRequest,
|
||||
description = "Data needed to create a Scheduled Message",
|
||||
content_type = "application/json"
|
||||
),
|
||||
responses(
|
||||
(status = 201, description = "Scheduled Message created", body = super::response::ScheduleMessageResponse),
|
||||
(status = 400, description = "Error", body = super::response::ScheduleMessageResponse),
|
||||
(status = 500, description = "Error creating Scheduled Message", body = super::response::ScheduleMessageResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn schedule_message(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::Json(payload): axum::Json<super::request::ScheduleMessageRequest>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::ScheduleMessageResponse>,
|
||||
) {
|
||||
let mut response = super::response::ScheduleMessageResponse::default();
|
||||
|
||||
if payload.is_valid() {
|
||||
if payload.status != textsender_models::message::scheduling::PENDING {
|
||||
response.message = String::from("scheduled message must be pending");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
} else {
|
||||
if can_schedule(&payload.scheduled) {
|
||||
let mut scheduled_message =
|
||||
textsender_models::message::scheduling::ScheduledMessage {
|
||||
scheduled: payload.scheduled,
|
||||
status: payload.status,
|
||||
..Default::default()
|
||||
};
|
||||
match scheduling_repo::insert(&pool, &scheduled_message, &payload.user_id).await
|
||||
{
|
||||
Ok((id, created)) => {
|
||||
scheduled_message.id = id;
|
||||
scheduled_message.created = Some(created);
|
||||
scheduled_message.user_id = payload.user_id;
|
||||
response.message = String::from("Message scheduled");
|
||||
response.data.push(scheduled_message);
|
||||
(axum::http::StatusCode::CREATED, axum::Json(response))
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Error inserting");
|
||||
(
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
response.message = String::from("Cannot schedule message");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
} else {
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
|
||||
/// Endpoint to get Scheduled Message
|
||||
#[utoipa::path(
|
||||
get,
|
||||
path = crate::caller::endpoints::GET_SCHEDULE_MESSAGE,
|
||||
params(
|
||||
("id" = uuid::Uuid, Path, description = "Id of Scheduled Message"),
|
||||
("user_id" = uuid::Uuid, Path, description = "User Id associated with the Scheduled Message")
|
||||
),
|
||||
responses(
|
||||
(status = 200, description = "Scheduled Message found", body = super::response::GetScheduledMessageResponse),
|
||||
(status = 400, description = "Error getting Scheduled Message", body = super::response::GetScheduledMessageResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn get_scheduled_messages(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::extract::Query(params): axum::extract::Query<
|
||||
super::request::GetScheduledMessageParams,
|
||||
>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::GetScheduledMessageResponse>,
|
||||
) {
|
||||
let mut response = super::response::GetScheduledMessageResponse::default();
|
||||
|
||||
let scheduled_messages = match params.id {
|
||||
Some(id) => match scheduling_repo::get(&pool, &id).await {
|
||||
Ok(scheduled_message) => {
|
||||
vec![scheduled_message]
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
},
|
||||
None => match params.user_id {
|
||||
Some(user_id) => match scheduling_repo::get_with_user_id(&pool, &user_id).await {
|
||||
Ok(scheduled_messages) => scheduled_messages,
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
},
|
||||
None => {
|
||||
response.message = String::from("Invalid parameter");
|
||||
return (axum::http::StatusCode::BAD_REQUEST, axum::Json(response));
|
||||
}
|
||||
},
|
||||
};
|
||||
|
||||
response.data = scheduled_messages;
|
||||
response.message = String::from(super::super::super::response::SUCCESSFUL);
|
||||
|
||||
(axum::http::StatusCode::OK, axum::Json(response))
|
||||
}
|
||||
|
||||
fn can_schedule(scheduled_time: &Option<time::OffsetDateTime>) -> bool {
|
||||
scheduled_time
|
||||
.is_some_and(|t| t - time::OffsetDateTime::now_utc() >= time::Duration::minutes(10))
|
||||
}
|
||||
|
||||
/// Endpoint to fetch Scheduled Message that is ready
|
||||
#[utoipa::path(
|
||||
get,
|
||||
path = crate::caller::endpoints::FETCH_SCHEDULED_MESSAGE,
|
||||
responses(
|
||||
(status = 200, description = "Scheduled Message found", body = super::response::FetchScheduledMessageResponse),
|
||||
(status = 400, description = "Error getting Scheduled Message", body = super::response::FetchScheduledMessageResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn fetch_scheduled_message(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::FetchScheduledMessageResponse>,
|
||||
) {
|
||||
let mut response = super::response::FetchScheduledMessageResponse::default();
|
||||
|
||||
match scheduling_repo::fetch(&pool).await {
|
||||
Ok(scheduled_message) => {
|
||||
response.data.push(scheduled_message);
|
||||
response.message = String::from(super::super::super::response::SUCCESSFUL);
|
||||
|
||||
(axum::http::StatusCode::OK, axum::Json(response))
|
||||
}
|
||||
Err(err) => match err {
|
||||
sqlx::Error::RowNotFound => {
|
||||
response.message = String::from("Nothing to fetch");
|
||||
|
||||
(axum::http::StatusCode::NO_CONTENT, axum::Json(response))
|
||||
}
|
||||
_ => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Error");
|
||||
(
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
)
|
||||
}
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/// Endpoint to update status of Scheduled Message
|
||||
#[utoipa::path(
|
||||
patch,
|
||||
path = crate::caller::endpoints::UPDATE_SCHEDULED_MESSAGE_STATUS,
|
||||
request_body(
|
||||
content = super::request::UpdateScheduledMessageStatusRequest,
|
||||
description = "Data needed to update status of Scheduled Message",
|
||||
content_type = "application/json"
|
||||
),
|
||||
responses(
|
||||
(status = 200, description = "Status updated", body = super::response::UpdateScheduledMessageStatusResponse),
|
||||
(status = 400, description = "Error", body = super::response::UpdateScheduledMessageStatusResponse),
|
||||
(status = 500, description = "Error updating", body = super::response::UpdateScheduledMessageStatusResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn update_scheduled_message_status(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::Json(payload): axum::Json<super::request::UpdateScheduledMessageStatusRequest>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::UpdateScheduledMessageStatusResponse>,
|
||||
) {
|
||||
let mut response = super::response::UpdateScheduledMessageStatusResponse::default();
|
||||
|
||||
if payload.is_valid() {
|
||||
match scheduling_repo::get(&pool, &payload.scheduled_message_id).await {
|
||||
Ok(scheduled_message) => {
|
||||
let old_status = scheduled_message.status;
|
||||
if old_status == payload.status {
|
||||
response.message = String::from("No change");
|
||||
(axum::http::StatusCode::NOT_MODIFIED, axum::Json(response))
|
||||
} else {
|
||||
match scheduling_repo::update_status(
|
||||
&pool,
|
||||
&scheduled_message.id,
|
||||
&payload.status,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => {
|
||||
response.message =
|
||||
String::from(super::super::super::response::SUCCESSFUL);
|
||||
response.data = vec![super::response::ChangedStatus { old_status }];
|
||||
(axum::http::StatusCode::OK, axum::Json(response))
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Invalid");
|
||||
(
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Invalid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
} else {
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
|
||||
/// Endpoint to create a Scheudled Message Event
|
||||
#[utoipa::path(
|
||||
post,
|
||||
path = crate::caller::endpoints::CREATE_SCHEDULED_MESSAGE_EVENT,
|
||||
request_body(
|
||||
content = super::request::CreateScheduleMessageEventRequest,
|
||||
description = "Data needed to create a Scheduled Message Event",
|
||||
content_type = "application/json"
|
||||
),
|
||||
responses(
|
||||
(status = 201, description = "Scheduled Message Event created", body = super::response::CreateScheduleMessageEventResponse),
|
||||
(status = 400, description = "Error", body = super::response::CreateScheduleMessageEventResponse),
|
||||
(status = 500, description = "Error creating Scheduled Message Event", body = super::response::CreateScheduleMessageEventResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn create_scheduled_message_event(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::Json(payload): axum::Json<super::request::CreateScheduleMessageEventRequest>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::CreateScheduleMessageEventResponse>,
|
||||
) {
|
||||
let mut response = super::response::CreateScheduleMessageEventResponse::default();
|
||||
|
||||
if payload.is_valid() {
|
||||
match scheduling_repo::get(&pool, &payload.scheduled_message_id).await {
|
||||
Ok(_) => {}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Scheduled Message does not exist");
|
||||
return (axum::http::StatusCode::BAD_REQUEST, axum::Json(response));
|
||||
}
|
||||
}
|
||||
|
||||
match contact_repo::get(&pool, &payload.contact_id).await {
|
||||
Ok(_) => {}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Contact does not exist");
|
||||
return (axum::http::StatusCode::BAD_REQUEST, axum::Json(response));
|
||||
}
|
||||
}
|
||||
|
||||
match message_repo::get(&pool, &payload.message_id).await {
|
||||
Ok(_) => {}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Contact does not exist");
|
||||
return (axum::http::StatusCode::BAD_REQUEST, axum::Json(response));
|
||||
}
|
||||
}
|
||||
|
||||
let mut scme = textsender_models::message::scheduling::ScheduledMessageEvent {
|
||||
contact_id: payload.contact_id,
|
||||
message_id: payload.message_id,
|
||||
scheduled_message_id: payload.scheduled_message_id,
|
||||
..Default::default()
|
||||
};
|
||||
match scheduling_repo::sched_msg_event::insert(&pool, &scme).await {
|
||||
Ok((id, created)) => {
|
||||
scme.id = id;
|
||||
scme.created = Some(created);
|
||||
response.message = String::from(super::super::super::response::SUCCESSFUL);
|
||||
response.data.push(scme);
|
||||
(axum::http::StatusCode::CREATED, axum::Json(response))
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Error inserting");
|
||||
(
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
response.message = String::from("Request body is not valid");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
|
||||
/// Endpoint to get Scheduled Message Event
|
||||
#[utoipa::path(
|
||||
get,
|
||||
path = crate::caller::endpoints::GET_SCHEDULED_MESSAGE_EVENT,
|
||||
params(
|
||||
("id" = uuid::Uuid, Path, description = "Id of Scheduled Message Event"),
|
||||
("scheduled_message_id" = uuid::Uuid, Path, description = "Scheduled Message Id associated with the Scheduled Message Event")
|
||||
),
|
||||
responses(
|
||||
(status = 200, description = "Scheduled Message Event found", body = super::response::GetScheduledMessageEventResponse),
|
||||
(status = 400, description = "Invalid request", body = super::response::GetScheduledMessageEventResponse),
|
||||
(status = 500, description = "Error getting Scheduled Message Event", body = super::response::GetScheduledMessageEventResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn get_scheduled_message_events(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::extract::Query(params): axum::extract::Query<
|
||||
super::request::GetScheduledMessageEventParams,
|
||||
>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::GetScheduledMessageEventResponse>,
|
||||
) {
|
||||
let mut response = super::response::GetScheduledMessageEventResponse::default();
|
||||
|
||||
let scheduled_message_events = match params.id {
|
||||
Some(id) => match scheduling_repo::sched_msg_event::get(&pool, &id).await {
|
||||
Ok(scheduled_message) => {
|
||||
vec![scheduled_message]
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
},
|
||||
None => match params.scheduled_message_id {
|
||||
Some(scheduled_message_id) => {
|
||||
match scheduling_repo::sched_msg_event::get_with_scheduled_message_id(
|
||||
&pool,
|
||||
&scheduled_message_id,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(scheduled_messages) => scheduled_messages,
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return (
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
None => {
|
||||
response.message = String::from("Invalid parameter");
|
||||
return (axum::http::StatusCode::BAD_REQUEST, axum::Json(response));
|
||||
}
|
||||
},
|
||||
};
|
||||
|
||||
response.data = scheduled_message_events;
|
||||
response.message = String::from(super::super::super::response::SUCCESSFUL);
|
||||
|
||||
(axum::http::StatusCode::OK, axum::Json(response))
|
||||
}
|
||||
|
||||
/// Endpoint to delete the Scheduled Message Event
|
||||
#[utoipa::path(
|
||||
delete,
|
||||
path = crate::caller::endpoints::DELETE_SCHEDULED_MESSAGE_EVENT,
|
||||
params(("id" = uuid::Uuid, Path, description = "Scheduled Message Event Id")),
|
||||
responses(
|
||||
(status = 200, description = "Deleted", body = super::response::DeleteScheduledMessageEventResponse),
|
||||
(status = 400, description = "Bad request", body = super::response::DeleteScheduledMessageEventResponse),
|
||||
(status = 500, description = "Error deleting", body = super::response::DeleteScheduledMessageEventResponse)
|
||||
)
|
||||
)]
|
||||
pub async fn delete_scheduled_message_event(
|
||||
axum::Extension(pool): axum::Extension<sqlx::PgPool>,
|
||||
axum::extract::Path(id): axum::extract::Path<uuid::Uuid>,
|
||||
) -> (
|
||||
axum::http::StatusCode,
|
||||
axum::Json<super::response::DeleteScheduledMessageEventResponse>,
|
||||
) {
|
||||
let mut response = super::response::DeleteScheduledMessageEventResponse::default();
|
||||
match scheduling_repo::sched_msg_event::get(&pool, &id).await {
|
||||
Ok(scheduled_message_event) => {
|
||||
match scheduling_repo::sched_msg_event::delete(&pool, &scheduled_message_event)
|
||||
.await
|
||||
{
|
||||
Ok(_) => {
|
||||
response.message = String::from(super::super::super::response::SUCCESSFUL);
|
||||
response.data.push(scheduled_message_event);
|
||||
(axum::http::StatusCode::OK, axum::Json(response))
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Error deleting");
|
||||
(
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
axum::Json(response),
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
response.message = String::from("Bad request");
|
||||
(axum::http::StatusCode::BAD_REQUEST, axum::Json(response))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,6 @@
|
||||
pub mod contact;
|
||||
pub mod instant;
|
||||
pub mod message;
|
||||
|
||||
pub mod endpoints {
|
||||
pub const ROOT: &str = "/";
|
||||
@@ -9,6 +11,31 @@ pub mod endpoints {
|
||||
pub const GET_CONTACT: &str = "/api/v1/contact";
|
||||
/// Constant for updating names of a Contact endpoint
|
||||
pub const UPDATE_CONTACT_NAME: &str = "/api/v1/contact/update";
|
||||
/// Constant for adding message endpoint
|
||||
pub const ADD_MESSAGE: &str = "/api/v1/message/new";
|
||||
/// Constant for getting messages endpoint
|
||||
pub const GET_MESSAGE: &str = "/api/v1/message";
|
||||
/// Constant for scheduling message endpoint
|
||||
pub const SCHEDULE_MESSAGE: &str = "/api/v1/schedule/message";
|
||||
/// Constant for getting scheduled message endpoint
|
||||
pub const GET_SCHEDULE_MESSAGE: &str = "/api/v1/schedule/message";
|
||||
/// Constant for updating scheduled message status endpoint
|
||||
pub const UPDATE_SCHEDULED_MESSAGE_STATUS: &str = "/api/v1/schedule/message/status/update";
|
||||
/// Constant for creating Scheduled Message Event endpoint
|
||||
pub const CREATE_SCHEDULED_MESSAGE_EVENT: &str = "/api/v1/schedule/message/event";
|
||||
/// Constant for getting Scheduled Message Event endpoint
|
||||
pub const GET_SCHEDULED_MESSAGE_EVENT: &str = "/api/v1/schedule/message/event";
|
||||
/// Constant for deleting Scheduled Message Event endpoint
|
||||
pub const DELETE_SCHEDULED_MESSAGE_EVENT: &str = "/api/v1/schedule/message/event/{id}";
|
||||
/// Constant for fetching Scheduled Message endpoint
|
||||
pub const FETCH_SCHEDULED_MESSAGE: &str = "/api/v1/schedule/message/fetch";
|
||||
/// Constant for recording Message Event Response endpoint
|
||||
pub const RECORD_MESSAGE_EVENT_RESPONSE: &str =
|
||||
"/api/v1/schedule/message/event/response/record";
|
||||
/// Constant for getting Message Event Response endpoint
|
||||
pub const GET_MESSAGE_EVENT_RESPONSE: &str = "/api/v1/schedule/message/event/response";
|
||||
/// Constant for sending a message instantly. No scheduling
|
||||
pub const INSTANT_MESSAGE: &str = "/api/v1/instant/message";
|
||||
}
|
||||
|
||||
pub mod response {
|
||||
|
||||
+92
-4
@@ -10,13 +10,23 @@ pub mod host {
|
||||
pub mod init {
|
||||
use std::time::Duration;
|
||||
|
||||
use axum::routing::{get, patch, post};
|
||||
use axum::routing::{delete, get, patch, post};
|
||||
use tower_http::timeout::TimeoutLayer;
|
||||
use utoipa::OpenApi;
|
||||
|
||||
use crate::caller::contact as contact_caller;
|
||||
use crate::caller::message as message_caller;
|
||||
use crate::caller::message::event as event_caller;
|
||||
use crate::caller::message::scheduling as scheduling_caller;
|
||||
|
||||
use contact_caller::endpoint as contact_endpoints;
|
||||
// use contact_caller::response as contact_responses;
|
||||
use contact_caller::response as contact_responses;
|
||||
use event_caller::endpoint as event_endpoints;
|
||||
use event_caller::response as event_responses;
|
||||
use message_caller::endpoint as message_endpoints;
|
||||
use message_caller::response as message_responses;
|
||||
use scheduling_caller::endpoint as scheduling_endpoints;
|
||||
use scheduling_caller::response as scheduling_responses;
|
||||
|
||||
mod cors {
|
||||
pub async fn configure_cors() -> tower_http::cors::CorsLayer {
|
||||
@@ -73,8 +83,14 @@ pub mod init {
|
||||
|
||||
#[derive(utoipa::OpenApi)]
|
||||
#[openapi(
|
||||
paths(contact_endpoints::create_contact,),
|
||||
components(schemas(contact_caller::response::AddContactResponse)),
|
||||
paths(contact_endpoints::create_contact,
|
||||
message_endpoints::create_message,
|
||||
scheduling_endpoints::schedule_message,
|
||||
),
|
||||
components(schemas(contact_responses::AddContactResponse, message_responses::AddMessageResponse,
|
||||
scheduling_responses::ScheduleMessageResponse,
|
||||
event_responses::RecordMessageEventResponse
|
||||
)),
|
||||
tags(
|
||||
(name = "textsender API", description = "Web API to manage texting")
|
||||
)
|
||||
@@ -101,6 +117,78 @@ pub mod init {
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::ADD_MESSAGE,
|
||||
post(message_endpoints::create_message).route_layer(axum::middleware::from_fn(
|
||||
crate::auth::auth::<axum::body::Body>,
|
||||
)),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::GET_MESSAGE,
|
||||
get(message_endpoints::get_messages).route_layer(axum::middleware::from_fn(
|
||||
crate::auth::auth::<axum::body::Body>,
|
||||
)),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::SCHEDULE_MESSAGE,
|
||||
post(scheduling_endpoints::schedule_message).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::GET_SCHEDULE_MESSAGE,
|
||||
get(scheduling_endpoints::get_scheduled_messages).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::CREATE_SCHEDULED_MESSAGE_EVENT,
|
||||
post(scheduling_endpoints::create_scheduled_message_event).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::GET_SCHEDULED_MESSAGE_EVENT,
|
||||
get(scheduling_endpoints::get_scheduled_message_events).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::DELETE_SCHEDULED_MESSAGE_EVENT,
|
||||
delete(scheduling_endpoints::delete_scheduled_message_event).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::UPDATE_SCHEDULED_MESSAGE_STATUS,
|
||||
patch(scheduling_endpoints::update_scheduled_message_status).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::FETCH_SCHEDULED_MESSAGE,
|
||||
get(scheduling_endpoints::fetch_scheduled_message).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::RECORD_MESSAGE_EVENT_RESPONSE,
|
||||
post(event_endpoints::record_message_event_response).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::GET_MESSAGE_EVENT_RESPONSE,
|
||||
get(event_endpoints::get_record_message_event_response).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.route(
|
||||
crate::caller::endpoints::INSTANT_MESSAGE,
|
||||
post(crate::caller::instant::endpoint::send_message).route_layer(
|
||||
axum::middleware::from_fn(crate::auth::auth::<axum::body::Body>),
|
||||
),
|
||||
)
|
||||
.layer(cors::configure_cors().await)
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
use sqlx::Row;
|
||||
|
||||
pub async fn insert(
|
||||
pool: &sqlx::PgPool,
|
||||
event: &textsender_models::message::event::MessageEventResponse,
|
||||
) -> Result<uuid::Uuid, sqlx::Error> {
|
||||
match sqlx::query(
|
||||
r#"
|
||||
INSERT INTO "message_event_responses" (scheduled_message_event_id, response, user_id, sent, contact_id, message_id, status)
|
||||
VALUES($1, $2, $3, $4, $5, $6, $7) RETURNING id;
|
||||
"#,
|
||||
)
|
||||
.bind(event.scheduled_message_event_id)
|
||||
.bind(&event.response)
|
||||
.bind(event.user_id)
|
||||
.bind(event.sent)
|
||||
.bind(event.contact_id)
|
||||
.bind(event.message_id)
|
||||
.bind(&event.status)
|
||||
.fetch_one(pool)
|
||||
.await
|
||||
{
|
||||
Ok(row) => {
|
||||
let id: uuid::Uuid = row.try_get("id")?;
|
||||
Ok(id)
|
||||
}
|
||||
Err(_) => Err(sqlx::Error::RowNotFound),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn get(
|
||||
pool: &sqlx::PgPool,
|
||||
id: &uuid::Uuid,
|
||||
) -> Result<textsender_models::message::event::MessageEventResponse, sqlx::Error> {
|
||||
match sqlx::query(
|
||||
r#"
|
||||
SELECT id, scheduled_message_event_id, response, user_id, sent, contact_id, message_id, status FROM "message_event_responses"
|
||||
WHERE
|
||||
id = $1;
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.fetch_one(pool)
|
||||
.await
|
||||
{
|
||||
Ok(row) => match parse_row(&row).await {
|
||||
Ok(event) => Ok(event),
|
||||
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::event::MessageEventResponse>, sqlx::Error> {
|
||||
match sqlx::query(
|
||||
r#"
|
||||
SELECT id, scheduled_message_event_id, response, user_id, sent, contact_id, message_id, status FROM "message_event_responses"
|
||||
WHERE
|
||||
user_id = $1;
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.fetch_all(pool)
|
||||
.await
|
||||
{
|
||||
Ok(rows) => {
|
||||
let mut events: Vec<textsender_models::message::event::MessageEventResponse> =
|
||||
Vec::new();
|
||||
for row in rows {
|
||||
match parse_row(&row).await {
|
||||
Ok(event) => {
|
||||
events.push(event);
|
||||
}
|
||||
Err(err) => {
|
||||
return Err(err);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(events)
|
||||
}
|
||||
Err(_) => Err(sqlx::Error::RowNotFound),
|
||||
}
|
||||
}
|
||||
|
||||
async fn parse_row(
|
||||
row: &sqlx::postgres::PgRow,
|
||||
) -> Result<textsender_models::message::event::MessageEventResponse, sqlx::Error> {
|
||||
println!("Parsing MER");
|
||||
let id: uuid::Uuid = row.try_get("id")?;
|
||||
let response: serde_json::Value = row.try_get("response")?;
|
||||
let status: String = row.try_get("status")?;
|
||||
let contact_id: uuid::Uuid = row.try_get("contact_id")?;
|
||||
let message_id: uuid::Uuid = row.try_get("message_id")?;
|
||||
let sent: time::OffsetDateTime = row.try_get("sent")?;
|
||||
let scheduled_message_event_id: uuid::Uuid = row.try_get("scheduled_message_event_id")?;
|
||||
let user_id: uuid::Uuid = row.try_get("user_id")?;
|
||||
|
||||
Ok(textsender_models::message::event::MessageEventResponse {
|
||||
id,
|
||||
response,
|
||||
contact_id,
|
||||
message_id,
|
||||
status,
|
||||
sent: Some(sent),
|
||||
scheduled_message_event_id: Some(scheduled_message_event_id),
|
||||
user_id,
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,94 @@
|
||||
pub mod event;
|
||||
pub mod scheduling;
|
||||
|
||||
use sqlx::Row;
|
||||
|
||||
pub async fn insert(
|
||||
pool: &sqlx::PgPool,
|
||||
message: &textsender_models::message::Message,
|
||||
user_id: &uuid::Uuid,
|
||||
) -> Result<uuid::Uuid, sqlx::Error> {
|
||||
match sqlx::query(
|
||||
r#"
|
||||
INSERT INTO "messages" (content, user_id)
|
||||
VALUES($1, $2) RETURNING id;
|
||||
"#,
|
||||
)
|
||||
.bind(&message.content)
|
||||
.bind(user_id)
|
||||
.fetch_one(pool)
|
||||
.await
|
||||
{
|
||||
Ok(row) => {
|
||||
let id: uuid::Uuid = row.try_get("id")?;
|
||||
Ok(id)
|
||||
}
|
||||
Err(_) => Err(sqlx::Error::RowNotFound),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn get(
|
||||
pool: &sqlx::PgPool,
|
||||
id: &uuid::Uuid,
|
||||
) -> Result<textsender_models::message::Message, sqlx::Error> {
|
||||
match sqlx::query(
|
||||
r#"
|
||||
SELECT id, content, user_id FROM "messages"
|
||||
WHERE
|
||||
id = $1;
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.fetch_one(pool)
|
||||
.await
|
||||
{
|
||||
Ok(row) => {
|
||||
let id: uuid::Uuid = row.try_get("id")?;
|
||||
|
||||
let content: String = row.try_get("content")?;
|
||||
let user_id: uuid::Uuid = row.try_get("user_id")?;
|
||||
Ok(textsender_models::message::Message {
|
||||
id: Some(id),
|
||||
content,
|
||||
user_id: Some(user_id),
|
||||
})
|
||||
}
|
||||
Err(_) => Err(sqlx::Error::RowNotFound),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn get_with_user_id(
|
||||
pool: &sqlx::PgPool,
|
||||
user_id: &uuid::Uuid,
|
||||
) -> Result<Vec<textsender_models::message::Message>, sqlx::Error> {
|
||||
match sqlx::query(
|
||||
r#"
|
||||
SELECT id, content, user_id FROM "messages"
|
||||
WHERE
|
||||
user_id = $1;
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.fetch_all(pool)
|
||||
.await
|
||||
{
|
||||
Ok(rows) => {
|
||||
let mut messages: Vec<textsender_models::message::Message> = Vec::new();
|
||||
for row in rows {
|
||||
let id: uuid::Uuid = row.try_get("id")?;
|
||||
let content: String = row.try_get("content")?;
|
||||
let user_id: uuid::Uuid = row.try_get("user_id")?;
|
||||
|
||||
let message = textsender_models::message::Message {
|
||||
id: Some(id),
|
||||
content,
|
||||
user_id: Some(user_id),
|
||||
};
|
||||
messages.push(message);
|
||||
}
|
||||
|
||||
Ok(messages)
|
||||
}
|
||||
Err(_) => Err(sqlx::Error::RowNotFound),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,280 @@
|
||||
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),
|
||||
},
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -1 +1,2 @@
|
||||
pub mod contact;
|
||||
pub mod message;
|
||||
|
||||
+1336
-110
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user