A lot of changes #12

Merged
phoenix merged 20 commits from alot_of_changes into main 2026-06-27 17:56:23 -04:00
21 changed files with 4380 additions and 488 deletions
+1 -1
View File
@@ -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
View File
@@ -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"
+91
View File
@@ -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
+1 -1
View File
@@ -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
View File
File diff suppressed because it is too large Load Diff
+4 -7
View File
@@ -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
View File
@@ -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
+33
View File
@@ -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'))
);
-49
View File
@@ -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'))
-- );
+1 -1
View File
@@ -171,7 +171,7 @@ pub mod endpoint {
}
}
// Endpoint to get songs
/// Endpoint to get Contacts
#[utoipa::path(
get,
path = super::super::endpoints::GET_CONTACT,
+171
View File
@@ -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))
}
}
}
+167
View File
@@ -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))
}
}
+154
View File
@@ -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))
}
}
+549
View File
@@ -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))
}
}
}
}
+27
View File
@@ -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
View File
@@ -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)
}
+112
View File
@@ -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,
})
}
+94
View File
@@ -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),
}
}
+280
View File
@@ -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
View File
@@ -1 +1,2 @@
pub mod contact;
pub mod message;
+1336 -110
View File
File diff suppressed because it is too large Load Diff