16 Commits
Author SHA1 Message Date
phoenix 71ba4a9c56 Update rust (#14)
catapult PR / Rustfmt (pull_request) Successful in 49s
catapult PR / Clippy (pull_request) Successful in 1m35s
catapult PR / Check (pull_request) Successful in 2m1s
Reviewed-on: #14
2026-07-06 22:53:39 -04:00
phoenix 8fc4bb7626 Update dependencies (#13)
catapult PR / Rustfmt (pull_request) Successful in 1m4s
catapult PR / Clippy (pull_request) Successful in 1m32s
catapult PR / Check (pull_request) Successful in 2m13s
Reviewed-on: #13
2026-07-04 21:26:22 -04:00
phoenix ec690cd23e v0.3.0 (#2)
Reviewed-on: #2
2026-06-27 17:37:22 -04:00
phoenix 2880e4142f Updated docker (#1)
Go / build (push) Successful in 39s
Go / build (pull_request) Successful in 1m10s
Reviewed-on: #1
2026-05-29 18:50:00 -04:00
phoenixandphoenix 0c2d638adf update go (#23)
Reviewed-on: #23
Co-authored-by: phoenix <mail@kundeng.us>
Co-committed-by: phoenix <mail@kundeng.us>
2026-05-03 16:55:25 -04:00
phoenixandphoenix 0691ce0adb Update go (#22)
Reviewed-on: #22
Co-authored-by: phoenix <mail@kundeng.us>
Co-committed-by: phoenix <mail@kundeng.us>
2026-04-04 22:21:25 -04:00
phoenixandphoenix 40d92ab7d3 tsk-20: Modify how Twilio message is sent out (#21)
Closes #20

Reviewed-on: #21
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2026-01-01 20:58:03 +00:00
phoenixandphoenix 7ec1f58f6a message_event_response-changes (#19)
Reviewed-on: #19
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-12-22 22:35:48 +00:00
phoenixandphoenix 66c7fc261f tsk-11: Obtain refresh token (#18)
Closes #11

Reviewed-on: #18
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-12-04 03:13:13 +00:00
phoenixandphoenix a4fefd5bc6 tsk-10: Dockerize app (#17)
Closes #10

Reviewed-on: #17
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-12-02 22:02:33 +00:00
phoenixandphoenix c809815549 tsk-14: Record when message was sent (#16)
Closes #14

Reviewed-on: #16
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-12-01 21:24:05 +00:00
phoenixandphoenix 04da7785c8 tsk-8: Schedule message (#15)
Closes #8

Reviewed-on: #15
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-29 19:28:05 +00:00
phoenixandphoenix 9af11dc961 tsk-7: Verify if scheduled message can be sent (#13)
Closes #7

Reviewed-on: #13
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-28 18:23:47 +00:00
phoenixandphoenix 6fa19ea4f6 tsk-6: Check queue (#12)
Closes #6

Reviewed-on: #12
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-28 18:01:18 +00:00
phoenixandphoenix cef7f8a150 tsk-5: Fetch token (#9)
Closes #5

Reviewed-on: #9
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-28 16:54:11 +00:00
phoenixandphoenix c47941acaf tsk-3: Create async loop (#4)
Closes #3

Reviewed-on: #4
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-27 01:26:09 +00:00
28 changed files with 3871 additions and 116 deletions
+8
View File
@@ -0,0 +1,8 @@
target/
.git/
.gitea/
.github/
*.DS_Store
+8
View File
@@ -0,0 +1,8 @@
AUTH_URL=https://auth.txt.com
API_URL=https://txt.com
SERVICE_USERNAME=suave
SERVICE_PASSPHRASE=9238urc9328nr329
TWILIO_AUTH_SID=9M438C93R943U4329MCU43C34U
TWILIO_SERVICE_SID=9M4J3X8439U398NUVT3342MC349C348T
TWILIO_AUTH_TOKEN="OJ8mU98U8UUU098u0U08kd9IDKaeoijtritjerDFISFGOS"
TWILIO_PHONE_NUMBER=10123456789
+8
View File
@@ -0,0 +1,8 @@
AUTH_URL=https://auth.txt.com
API_URL=https://txt.com
SERVICE_USERNAME=suave
SERVICE_PASSPHRASE=9238urc9328nr329
TWILIO_AUTH_SID=9M438C93R943U4329MCU43C34U
TWILIO_SERVICE_SID=9M4J3X8439U398NUVT3342MC349C348T
TWILIO_AUTH_TOKEN=OJ8mU98U8UUU098u0U08kd9IDKaeoijtritjerDFISFGOS
TWILIO_PHONE_NUMBER=10123456789
+70
View File
@@ -0,0 +1,70 @@
name: catapult PR
on:
pull_request:
branches:
- main
concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: true
jobs:
check:
name: Check
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v6
- uses: actions-rust-lang/setup-rust-toolchain@v1
with:
toolchain: 1.96.1
- 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 check
fmt:
name: Rustfmt
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v6
- uses: actions-rust-lang/setup-rust-toolchain@v1
with:
toolchain: 1.96.1
- 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
clippy:
name: Clippy
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v6
- uses: actions-rust-lang/setup-rust-toolchain@v1
with:
toolchain: 1.96.1
- run: rustup component add clippy
- 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 clippy -- -D warnings
-44
View File
@@ -1,44 +0,0 @@
name: Go
on:
push:
branches: [ main ]
pull_request:
branches: [ main ]
jobs:
build:
runs-on: ubuntu-24.04 # You can change this to macos-latest or windows-latest if needed
steps:
- uses: actions/checkout@v5
- name: Set up Go
uses: actions/setup-go@v6
with:
go-version: '1.25.4' # You can specify a specific version or 'stable'
- name: Build
run: |
echo "Initializing config"
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
go env -w GOPRIVATE='${{ secrets.GIT_HOST_ROOT }}'
echo "Creating local .gitconfig"
touch ~/.gitconfig
cat > ~/.gitconfig << "EOF"
[url "ssh://git@${{ secrets.GIT_HOST_ROOT }}"]
insteadOf = https://${{ secrets.GIT_HOST_ROOT }}
EOF
make build
- name: Test
run: go test -v ./...
+4 -1
View File
@@ -1 +1,4 @@
catapult /target
.env
.env.docker
.env.local
Generated
+2823
View File
File diff suppressed because it is too large Load Diff
+16
View File
@@ -0,0 +1,16 @@
[package]
name = "catapult"
version = "0.3.2"
rust-version = "1.96.1"
edition = "2024"
[dependencies]
tokio = { version = "1.52.3", features = ["full"] }
reqwest = { version = "0.13.4", features = ["json", "stream", "multipart"] }
serde_json = { version = "1.0.150" }
time = { version = "0.3.53", features = ["formatting", "macros", "parsing", "serde"] }
uuid = { version = "1.23.3", features = ["v4", "serde"] }
textsender_models = { git = "ssh://git@git.kundeng.us/phoenix/textsender_models.git", tag = "v0.5.1" }
swoosh = { git = "ssh://git@git.kundeng.us/phoenix/swoosh.git", tag = "v0.5.2" }
[dev-dependencies]
+61
View File
@@ -0,0 +1,61 @@
# Stage 1: Build the application
FROM rust:1.96.1 as builder
# Set the working directory inside the container
WORKDIR /usr/src/app
# Install build dependencies if needed (e.g., git for cloning)
RUN apt-get update && apt-get install -y --no-install-recommends \
pkg-config libssl3 \
ca-certificates \
openssh-client git \
&& rm -rf /var/lib/apt/lists/*
# << --- ADD HOST KEY HERE --- >>
# Replace 'yourgithost.com' with the actual hostname (e.g., github.com)
RUN mkdir -p -m 0700 ~/.ssh && \
ssh-keyscan git.kundeng.us >> ~/.ssh/known_hosts
# Copy Cargo manifests
COPY Cargo.toml Cargo.lock ./
# Build *only* dependencies to leverage Docker cache
# This dummy build caches dependencies as a separate layer
RUN --mount=type=ssh mkdir src && \
echo "fn main() {println!(\"if you see this, the build broke\")}" > src/main.rs && \
cargo build --release --quiet && \
rm -rf src target/release/deps/catapult*
# Copy the actual source code
COPY src ./src
# If you have other directories like `templates` or `static`, copy them too
COPY .env ./.env
# << --- SSH MOUNT ADDED HERE --- >>
# Build *only* dependencies to leverage Docker cache
# This dummy build caches dependencies as a separate layer
# Mount the SSH agent socket for this command
RUN --mount=type=ssh \
cargo build --release --quiet
# Stage 2: Create the final, smaller runtime image
# Use a minimal base image like debian-slim or even distroless for security/size
FROM debian:trixie-slim
# Install runtime dependencies if needed (e.g., SSL certificates)
RUN apt-get update && apt-get install -y ca-certificates libssl-dev libssl3 && rm -rf /var/lib/apt/lists/*
# Set the working directory
WORKDIR /usr/local/bin
# Copy the compiled binary from the builder stage
# Replace 'catapult' with the actual name of your binary (usually the crate name)
COPY --from=builder /usr/src/app/target/release/catapult .
# Copy other necessary files like .env (if used for runtime config) or static assets
# It's generally better to configure via environment variables in Docker though
COPY --from=builder /usr/src/app/.env .
# Set the command to run your application
# Ensure this matches the binary name copied above
CMD ["./catapult"]
-21
View File
@@ -1,21 +0,0 @@
VERSION ?= $(shell git describe --tags 2>/dev/null || echo "dev")
COMMIT ?= $(shell git rev-parse --short HEAD)
BUILD_TIME ?= $(shell date -u +%Y-%m-%dT%H:%M:%SZ)
GO_VERSION ?= $(shell go version | awk '{print $$3}')
.PHONY: build
build:
go build -ldflags="\
-X 'git.kundeng.us/phoenix/catapult/version.Version=$(VERSION)' \
-X 'git.kundeng.us/phoenix/catapult/version.BuildTime=$(BUILD_TIME)' \
-X 'git.kundeng.us/phoenix/catapult/version.Commit=$(COMMIT)' \
-X 'git.kundeng.us/phoenix/catapult/version.GoVersion=$(GO_VERSION)'" \
-o catapult main.go
.PHONY: install
install:
go install -ldflags="\
-X 'git.kundeng.us/phoenix/catapult/version.Version=$(VERSION)' \
-X 'git.kundeng.us/phoenix/catapult/version.BuildTime=$(BUILD_TIME)' \
-X 'git.kundeng.us/phoenix/catapult/version.Commit=$(COMMIT)' \
-X 'git.kundeng.us/phoenix/catapult/version.GoVersion=$(GO_VERSION)'"
-8
View File
@@ -1,8 +0,0 @@
# Catapult
A service used to send text messages
Building the service
```
make build
```
+11
View File
@@ -0,0 +1,11 @@
version: '3.8' # Use a recent version
services:
songparser:
build:
context: .
ssh: ["default"] # Uses host's SSH agent
container_name: catapult
env_file:
- .env
restart: unless-stopped # Optional: Restart policy
-3
View File
@@ -1,3 +0,0 @@
module git.kundeng.us/phoenix/catapult
go 1.25.4
-17
View File
@@ -1,17 +0,0 @@
package version
import "fmt"
var (
Version = "dev"
BuildTime = "unknown"
Commit = "unknown"
GoVersion = "unknown"
)
func String() string {
return fmt.Sprintf(
"Version: %s\nBuild Date: %s\nCommit: %s\nGo Version: %s",
Version, BuildTime, Commit, GoVersion,
)
}
-22
View File
@@ -1,22 +0,0 @@
package main
import (
"flag"
"fmt"
"git.kundeng.us/phoenix/catapult/internal/version"
)
const App_Name = "catapult"
func main() {
fmt.Println(App_Name)
versionFlag := flag.Bool("version", false, "Print version information")
flag.Parse()
if *versionFlag {
fmt.Println(version.String())
return
}
}
+54
View File
@@ -0,0 +1,54 @@
/// Seconds to sleep after finising a task
pub const SECONDS_TO_SLEEP: u64 = 5;
pub const APP_NAME: &str = "catapult";
#[derive(Clone, Debug)]
pub struct App {
pub api_url: String,
pub auth_url: String,
pub service_username: String,
pub service_passphrase: String,
pub twilio_config: textsender_models::config::auxiliary::TwilioConfig,
}
pub fn load_app() -> Result<App, std::io::Error> {
use textsender_models::envy;
let auth_url_env = envy::environment::get_env("AUTH_URL");
let api_url_env = envy::environment::get_env("API_URL");
let service_username_env = envy::environment::get_env("SERVICE_USERNAME");
let service_passphrase_env = envy::environment::get_env("SERVICE_PASSPHRASE");
let envs = vec![
api_url_env,
auth_url_env,
service_username_env,
service_passphrase_env,
];
let check_envs = |vs: &Vec<envy::EnvVar>| -> (bool, Option<String>) {
for env in vs {
if env.value.is_empty() {
return (false, Some(format!("Key {} is not provided", env.key)));
}
}
(true, None)
};
let (envs_valid, reason) = check_envs(&envs);
if envs_valid {
match textsender_models::config::auxiliary::load_config() {
Ok(twilio_config) => Ok(App {
api_url: envs[0].value.clone(),
auth_url: envs[1].value.clone(),
service_username: envs[2].value.clone(),
service_passphrase: envs[3].value.clone(),
twilio_config,
}),
Err(err) => Err(err),
}
} else {
Err(std::io::Error::other(reason.unwrap()))
}
}
+3
View File
@@ -0,0 +1,3 @@
pub mod app;
pub mod service;
pub mod version;
+55
View File
@@ -0,0 +1,55 @@
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let args = std::env::args().collect();
if catapult::version::contains_version(&args) {
catapult::version::print_version();
std::process::exit(-1);
}
let app = match catapult::app::load_app() {
Ok(app) => app,
Err(err) => {
eprintln!("Error: {err:?}");
std::process::exit(-1);
}
};
println!("App: {app:?}");
let auth = catapult::service::auth::Auth { app: app.clone() };
let token = match auth.get_token().await {
Ok(token) => token,
Err(err) => {
eprintln!("Error getting token");
eprintln!("Error: {err:?}");
std::process::exit(-1);
}
};
let mut svc = catapult::service::core::Service {
app: app.clone(),
token,
};
loop {
if catapult::service::auth::has_token_expired(&svc.token).await {
match auth.get_refresh_token(&svc.token).await {
Ok(refresh) => {
println!("Refresh token: {refresh:?}");
svc.token = refresh;
}
Err(err) => {
eprintln!("Error: {err:?}");
std::process::exit(-1);
}
}
}
svc.do_the_work().await;
tokio::time::sleep(tokio::time::Duration::from_secs(
catapult::app::SECONDS_TO_SLEEP,
))
.await;
}
}
+116
View File
@@ -0,0 +1,116 @@
pub struct Auth {
pub app: crate::app::App,
}
impl Auth {
pub async fn get_token(&self) -> Result<textsender_models::token::Login, std::io::Error> {
let client = reqwest::Client::new();
let endpoint = String::from("api/v1/service/login");
let api_url = format!("{}/{endpoint}", self.app.auth_url);
println!("Url: {api_url:?}");
let payload = serde_json::json!({
"username": &self.app.service_username,
"passphrase": &self.app.service_passphrase,
});
println!("Payload: {payload:?}");
match client.post(api_url).json(&payload).send().await {
Ok(response) => match response.status() {
reqwest::StatusCode::OK => {
println!("Response: {response:?}");
match response.text().await {
Ok(val) => match parse_token_response(&val).await {
Ok(login_result) => Ok(login_result),
Err(err) => Err(err),
},
Err(err) => Err(std::io::Error::other(err)),
}
}
_ => Err(std::io::Error::other("Error getting token")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
pub async fn get_refresh_token(
&self,
login_result: &textsender_models::token::Login,
) -> Result<textsender_models::token::Login, std::io::Error> {
let client = reqwest::Client::new();
let endpoint = String::from("api/v1/token/refresh");
let api_url = format!("{}/{endpoint}", self.app.auth_url);
println!("Url: {api_url:?}");
let payload = serde_json::json!({
"access_token": login_result.access_token
});
match client.post(api_url).json(&payload).send().await {
Ok(response) => match response.status() {
reqwest::StatusCode::OK => {
println!("Response: {response:?}");
match response.text().await {
Ok(val) => match parse_token_response(&val).await {
Ok(login_result) => Ok(login_result),
Err(err) => Err(err),
},
Err(err) => Err(std::io::Error::other(err)),
}
}
_ => Err(std::io::Error::other("Error getting token")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
}
async fn parse_token_response(
response: &str,
) -> Result<textsender_models::token::LoginResult, std::io::Error> {
match serde_json::from_str::<serde_json::Value>(response) {
Ok(j) => {
println!("Good");
let message = &j["message"];
let data = &j["data"];
println!("Message: {message:?}");
println!("Data: {data:?}");
match j.get("data") {
Some(serde_json::Value::Array(lrs)) => {
if lrs.is_empty() {
Err(std::io::Error::other("Error response is empty"))
} else {
let lr = lrs[0].clone();
let ll: textsender_models::token::LoginResult =
match serde_json::from_value(lr) {
Ok(lr) => lr,
Err(err) => {
return Err(std::io::Error::other(err.to_string()));
}
};
Ok(ll)
}
}
_ => Err(std::io::Error::other("Error parsing response")),
}
}
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
pub async fn has_token_expired(token: &textsender_models::token::LoginResult) -> bool {
let now = time::OffsetDateTime::now_utc();
let expire = match time::OffsetDateTime::from_unix_timestamp(token.expires_in) {
Ok(res) => res,
Err(err) => {
eprintln!("Error: {err:?}");
return false;
}
};
now > expire
}
+64
View File
@@ -0,0 +1,64 @@
pub struct Contact {
pub app: crate::app::App,
pub token: textsender_models::token::LoginResult,
}
impl Contact {
pub async fn get(
&self,
contact_id: &uuid::Uuid,
) -> Result<Vec<textsender_models::contact::Contact>, std::io::Error> {
let client = reqwest::Client::new();
let endpoint = format!("api/v1/contact?id={contact_id}");
let api_url = format!("{}/{endpoint}", self.app.api_url);
println!("Url: {api_url:?}");
let (key, header) = super::auth_header(&self.token.access_token);
match client.get(api_url).header(key, header).send().await {
Ok(response) => match response.status() {
reqwest::StatusCode::OK => match response.text().await {
Ok(val) => match parse_response(&val).await {
Ok(login_result) => Ok(login_result),
Err(err) => Err(err),
},
Err(err) => Err(std::io::Error::other(err)),
},
_ => Err(std::io::Error::other("Error getting token")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
}
async fn parse_response(
response: &str,
) -> Result<Vec<textsender_models::contact::Contact>, std::io::Error> {
match serde_json::from_str::<serde_json::Value>(response) {
Ok(j) => match j.get("data") {
Some(serde_json::Value::Array(lrs)) => {
if lrs.is_empty() {
Err(std::io::Error::other("Error response is empty"))
} else {
let mut events: Vec<textsender_models::contact::Contact> = Vec::new();
for event in lrs.iter() {
let ll: textsender_models::contact::Contact =
match serde_json::from_value(event.clone()) {
Ok(lr) => lr,
Err(err) => {
return Err(std::io::Error::other(err.to_string()));
}
};
events.push(ll);
}
Ok(events)
}
}
_ => Err(std::io::Error::other("Error parsing response")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
+161
View File
@@ -0,0 +1,161 @@
pub struct Service {
pub app: crate::app::App,
pub token: textsender_models::token::Login,
}
impl Service {
pub async fn do_the_work(&mut self) {
println!("Checking queue");
use textsender_models::message::scheduling::ScheduledMessage;
let queue = super::queue::Queue {
app: self.app.clone(),
token: self.token.clone(),
};
let queue_item: Option<ScheduledMessage> = match queue.get_queue().await {
Ok(item) => Some(item),
Err(err) => {
eprintln!("Error: {err:?}");
None
}
};
if queue_item.is_none() {
println!("No queue item found");
return;
}
let queue_item = queue_item.unwrap();
println!("Queue item Id: {:?}", queue_item.id);
if !self.is_schedulable(&queue_item) {
println!("Not schedulable");
return;
}
println!("Can schedule");
let events_req = super::event::Event {
app: self.app.clone(),
token: self.token.clone(),
};
let scheduled_message_id = queue_item.id;
let events = match events_req.get(&scheduled_message_id).await {
Ok(events) => Some(events),
Err(err) => {
eprintln!("Error: {err:?}");
None
}
};
let contacts = super::contact::Contact {
app: self.app.clone(),
token: self.token.clone(),
};
let messages = super::message::Message {
app: self.app.clone(),
token: self.token.clone(),
};
let sch_msg = super::scheduler::ScheduledMessage {
app: self.app.clone(),
token: self.token.clone(),
};
let mer_req = super::mer::MessageEventResponse {
app: self.app.clone(),
token: self.token.clone(),
};
let events = events.unwrap_or_default();
let mut sendmsg = swoosh::SendMsg::default();
sendmsg.load_config(
&self.app.twilio_config.auth_token,
&self.app.twilio_config.phone_number,
&self.app.twilio_config.service_sid,
&self.app.twilio_config.account_sid,
);
for event in &events {
println!("Event: {event:?}");
let contact = &match contacts.get(&event.contact_id).await {
Ok(contact) => contact,
Err(err) => {
eprintln!("Error: {err:?}");
continue;
}
}[0];
let message = &match messages.get(&event.message_id).await {
Ok(messages) => messages,
Err(err) => {
eprintln!("Error: {err:?}");
continue;
}
}[0];
let scheduled_message = match sch_msg.get(&event.scheduled_message_id).await {
Ok(sch_msg_here) => sch_msg_here,
Err(err) => {
eprintln!("Error: {err:?}");
continue;
}
};
println!("Contact Id: {:?}", contact.id);
println!("Message Id: {:?}", message.id);
println!("Scheduled Message Id: {:?}", scheduled_message.id);
let param = swoosh::twilio::types::Parameters {
schedule: true,
schedule_at: scheduled_message.scheduled,
};
let twilio_config = &self.app.twilio_config;
println!("Config: {twilio_config:?}");
println!("SM: {scheduled_message:?}");
sendmsg.load_message(&message.content);
sendmsg.load_recipient(&contact.phone_number);
match sendmsg.send_message(&param).await {
Ok(resp) => {
println!("Message scheduled");
// MER response
let result = swoosh::twilio::api::response_to_json(resp).await;
let mer = textsender_models::message::event::MessageEventResponse {
contact_id: contact.id.unwrap(),
scheduled_message_event_id: Some(event.id),
response: result,
message_id: message.id.unwrap(),
sent: scheduled_message.scheduled,
user_id: scheduled_message.user_id,
..Default::default()
};
match mer_req.record_message_event_response(&mer).await {
Ok(response) => {
println!("Recording MER");
println!("Response: {response:?}");
}
Err(err) => {
eprintln!("Error recording mer");
eprintln!("Error: {err:?}");
}
}
}
Err(err) => {
eprintln!("Message not sent");
eprintln!("Error: {err:?}");
}
}
}
}
fn is_schedulable(
&self,
queue_item: &textsender_models::message::scheduling::ScheduledMessage,
) -> bool {
match queue_item.scheduled {
Some(date) => date > time::OffsetDateTime::now_utc(),
None => false,
}
}
}
+70
View File
@@ -0,0 +1,70 @@
pub struct Event {
pub app: crate::app::App,
pub token: textsender_models::token::LoginResult,
}
impl Event {
pub async fn get(
&self,
scheduled_message_id: &uuid::Uuid,
) -> Result<Vec<textsender_models::message::scheduling::ScheduledMessageEvent>, std::io::Error>
{
let client = reqwest::Client::new();
let endpoint =
format!("api/v1/schedule/message/event?scheduled_message_id={scheduled_message_id}");
let api_url = format!("{}/{endpoint}", self.app.api_url);
println!("Url: {api_url:?}");
let (key, header) = super::auth_header(&self.token.access_token);
match client.get(api_url).header(key, header).send().await {
Ok(response) => match response.status() {
reqwest::StatusCode::OK => match response.text().await {
Ok(val) => match parse_response(&val).await {
Ok(login_result) => Ok(login_result),
Err(err) => Err(err),
},
Err(err) => Err(std::io::Error::other(err)),
},
_ => Err(std::io::Error::other(
"Error getting Scheduled Message Event",
)),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
}
async fn parse_response(
response: &str,
) -> Result<Vec<textsender_models::message::scheduling::ScheduledMessageEvent>, std::io::Error> {
match serde_json::from_str::<serde_json::Value>(response) {
Ok(j) => match j.get("data") {
Some(serde_json::Value::Array(lrs)) => {
if lrs.is_empty() {
Err(std::io::Error::other("Error response is empty"))
} else {
let mut events: Vec<
textsender_models::message::scheduling::ScheduledMessageEvent,
> = Vec::new();
for event in lrs.iter() {
let ll: textsender_models::message::scheduling::ScheduledMessageEvent =
match serde_json::from_value(event.clone()) {
Ok(lr) => lr,
Err(err) => {
return Err(std::io::Error::other(err.to_string()));
}
};
events.push(ll);
}
Ok(events)
}
}
_ => Err(std::io::Error::other("Error parsing response")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
+112
View File
@@ -0,0 +1,112 @@
pub struct MessageEventResponse {
pub app: crate::app::App,
pub token: textsender_models::token::LoginResult,
}
impl MessageEventResponse {
pub async fn record_message_event_response(
&self,
mer: &textsender_models::message::event::MessageEventResponse,
) -> Result<textsender_models::message::event::MessageEventResponse, std::io::Error> {
let client = reqwest::Client::new();
let endpoint = String::from("api/v1/schedule/message/event/response/record");
let api_url = format!("{}/{endpoint}", self.app.api_url);
println!("Url: {api_url:?}");
use time::format_description::well_known::Iso8601;
let sent = match mer.sent.unwrap().format(&Iso8601::DEFAULT) {
Ok(converted) => converted,
Err(err) => {
return Err(std::io::Error::other(err.to_string()));
}
};
let payload = serde_json::json!({
"scheduled_message_event_id": mer.scheduled_message_event_id.unwrap(),
"response": mer.response,
"user_id": mer.user_id,
"sent": sent,
"status": textsender_models::message::event::MESSAGE_EVENT_RESPONSE_STATUS_SCHEDULED,
"contact_id": mer.contact_id,
"message_id": mer.message_id
});
println!("Payload: {payload:?}");
let (key, header) = super::auth_header(&self.token.access_token);
match client
.post(api_url)
.json(&payload)
.header(key, header)
.send()
.await
{
Ok(response) => match response.status() {
reqwest::StatusCode::CREATED | reqwest::StatusCode::OK => {
println!("Response: {response:?}");
match response.text().await {
Ok(val) => match parse_response(&val).await {
Ok(login_result) => Ok(login_result),
Err(err) => Err(err),
},
Err(err) => Err(std::io::Error::other(err)),
}
}
_ => Err(std::io::Error::other(
"Error recording Message Event Response",
)),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
/*
*
* {
"scheduled_message_event_id": "{{scheduled_message_event_id}}",
"response": {{text_response_raw}},
"user_id": "{{user_id}}",
"sent": "2026-06-20T17:10:00Z",
"status": "{{mer_status_scheduled}}",
"contact_id": "{{contact_id}}",
"message_id": "{{message_id}}"
}
*
*/
}
async fn parse_response(
response: &str,
) -> Result<textsender_models::message::event::MessageEventResponse, std::io::Error> {
match serde_json::from_str::<serde_json::Value>(response) {
Ok(j) => {
println!("Good");
let message = &j["message"];
let data = &j["data"];
println!("Message: {message:?}");
println!("Data: {data:?}");
match j.get("data") {
Some(serde_json::Value::Array(lrs)) => {
if lrs.is_empty() {
Err(std::io::Error::other("Error response is empty"))
} else {
let lr = lrs[0].clone();
let ll: textsender_models::message::event::MessageEventResponse =
match serde_json::from_value(lr) {
Ok(lr) => lr,
Err(err) => {
return Err(std::io::Error::other(err.to_string()));
}
};
Ok(ll)
}
}
_ => Err(std::io::Error::other("Error parsing response")),
}
}
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
+64
View File
@@ -0,0 +1,64 @@
pub struct Message {
pub app: crate::app::App,
pub token: textsender_models::token::LoginResult,
}
impl Message {
pub async fn get(
&self,
message_id: &uuid::Uuid,
) -> Result<Vec<textsender_models::message::Message>, std::io::Error> {
let client = reqwest::Client::new();
let endpoint = format!("api/v1/message?id={message_id}");
let api_url = format!("{}/{endpoint}", self.app.api_url);
println!("Url: {api_url:?}");
let (key, header) = super::auth_header(&self.token.access_token);
match client.get(api_url).header(key, header).send().await {
Ok(response) => match response.status() {
reqwest::StatusCode::OK => match response.text().await {
Ok(val) => match parse_response(&val).await {
Ok(login_result) => Ok(login_result),
Err(err) => Err(err),
},
Err(err) => Err(std::io::Error::other(err)),
},
_ => Err(std::io::Error::other("Error getting token")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
}
async fn parse_response(
response: &str,
) -> Result<Vec<textsender_models::message::Message>, std::io::Error> {
match serde_json::from_str::<serde_json::Value>(response) {
Ok(j) => match j.get("data") {
Some(serde_json::Value::Array(lrs)) => {
if lrs.is_empty() {
Err(std::io::Error::other("Error response is empty"))
} else {
let mut events: Vec<textsender_models::message::Message> = Vec::new();
for event in lrs.iter() {
let ll: textsender_models::message::Message =
match serde_json::from_value(event.clone()) {
Ok(lr) => lr,
Err(err) => {
return Err(std::io::Error::other(err.to_string()));
}
};
events.push(ll);
}
Ok(events)
}
}
_ => Err(std::io::Error::other("Error parsing response")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
+16
View File
@@ -0,0 +1,16 @@
pub mod auth;
pub mod contact;
pub mod core;
pub mod event;
pub mod mer;
pub mod message;
pub mod queue;
pub mod scheduler;
pub fn auth_header(
access_token: &str,
) -> (reqwest::header::HeaderName, reqwest::header::HeaderValue) {
let bearer = format!("Bearer {}", access_token);
let header_value = reqwest::header::HeaderValue::from_str(&bearer).unwrap();
(reqwest::header::AUTHORIZATION, header_value)
}
+68
View File
@@ -0,0 +1,68 @@
pub struct Queue {
pub app: crate::app::App,
pub token: textsender_models::token::LoginResult,
}
impl Queue {
pub async fn get_queue(
&self,
) -> Result<textsender_models::message::scheduling::ScheduledMessage, std::io::Error> {
let client = reqwest::Client::new();
let endpoint = String::from("api/v1/schedule/message/fetch");
let api_url = format!("{}/{endpoint}", self.app.api_url);
println!("Url: {api_url:?}");
let (key, header) = super::auth_header(&self.token.access_token);
match client.get(api_url).header(key, header).send().await {
Ok(response) => match response.status() {
reqwest::StatusCode::OK => {
println!("Response: {response:?}");
match response.text().await {
Ok(val) => match parse_queue_response(&val).await {
Ok(login_result) => Ok(login_result),
Err(err) => Err(err),
},
Err(err) => Err(std::io::Error::other(err)),
}
}
_ => Err(std::io::Error::other("Error getting token")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
}
async fn parse_queue_response(
response: &str,
) -> Result<textsender_models::message::scheduling::ScheduledMessage, std::io::Error> {
match serde_json::from_str::<serde_json::Value>(response) {
Ok(j) => {
let message = &j["message"];
let data = &j["data"];
println!("Message: {message:?}");
println!("Data: {data:?}");
match j.get("data") {
Some(serde_json::Value::Array(lrs)) => {
if lrs.is_empty() {
Err(std::io::Error::other("Error response is empty"))
} else {
let lr = lrs[0].clone();
let ll: textsender_models::message::scheduling::ScheduledMessage =
match serde_json::from_value(lr) {
Ok(lr) => lr,
Err(err) => {
return Err(std::io::Error::other(err.to_string()));
}
};
Ok(ll)
}
}
_ => Err(std::io::Error::other("Error parsing response")),
}
}
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
+59
View File
@@ -0,0 +1,59 @@
pub struct ScheduledMessage {
pub app: crate::app::App,
pub token: textsender_models::token::LoginResult,
}
impl ScheduledMessage {
pub async fn get(
&self,
scheduled_message_id: &uuid::Uuid,
) -> Result<textsender_models::message::scheduling::ScheduledMessage, std::io::Error> {
let client = reqwest::Client::new();
let endpoint = format!("api/v1/schedule/message?id={scheduled_message_id}");
let api_url = format!("{}/{endpoint}", self.app.api_url);
println!("Url: {api_url:?}");
let (key, header) = super::auth_header(&self.token.access_token);
match client.get(api_url).header(key, header).send().await {
Ok(response) => match response.status() {
reqwest::StatusCode::OK => match response.text().await {
Ok(val) => match parse_response(&val).await {
Ok(login_result) => Ok(login_result),
Err(err) => Err(err),
},
Err(err) => Err(std::io::Error::other(err)),
},
_ => Err(std::io::Error::other("Error getting Scheduled Message")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
}
async fn parse_response(
response: &str,
) -> Result<textsender_models::message::scheduling::ScheduledMessage, std::io::Error> {
match serde_json::from_str::<serde_json::Value>(response) {
Ok(j) => match j.get("data") {
Some(serde_json::Value::Array(lrs)) => {
if lrs.is_empty() {
Err(std::io::Error::other("Error response is empty"))
} else {
let ll: textsender_models::message::scheduling::ScheduledMessage =
match serde_json::from_value(lrs[0].clone()) {
Ok(lr) => lr,
Err(err) => {
return Err(std::io::Error::other(err.to_string()));
}
};
Ok(ll)
}
}
_ => Err(std::io::Error::other("Error parsing response")),
},
Err(err) => Err(std::io::Error::other(err.to_string())),
}
}
+20
View File
@@ -0,0 +1,20 @@
pub fn print_version() {
let name = env!("CARGO_PKG_NAME");
let version = env!("CARGO_PKG_VERSION");
println!("{name:?} {version:?}");
}
pub fn contains_version(args: &Vec<String>) -> bool {
let valid_flags = vec!["-v", "--version"];
for arg in args {
for valid_flag in &valid_flags {
if valid_flag == arg {
return true;
}
}
}
false
}