19 Commits
Author SHA1 Message Date
phoenix f9c098b92f Adding license (#18)
catapult PR / Rustfmt (push) Successful in 41s
catapult PR / Clippy (push) Successful in 59s
catapult PR / Check (push) Successful in 1m44s
Reviewed-on: #18
2026-07-14 11:47:05 -04:00
phoenix 6801b1e635 Dependency name change (#17)
Reviewed-on: #17
2026-07-13 16:33:51 -04:00
phoenix 50aa9089f0 Update rust (#16)
catapult PR / Rustfmt (pull_request) Successful in 48s
catapult PR / Clippy (pull_request) Successful in 1m55s
catapult PR / Check (pull_request) Successful in 2m21s
Reviewed-on: #16
2026-07-10 18:33:46 -04:00
phoenix 8dbfd88aae Refactoring (#15)
catapult PR / Check (pull_request) Successful in 59s
catapult PR / Rustfmt (pull_request) Successful in 1m25s
catapult PR / Clippy (pull_request) Successful in 1m18s
Reviewed-on: #15
2026-07-07 15:51:06 -04:00
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
30 changed files with 3917 additions and 226 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
+73
View File
@@ -0,0 +1,73 @@
name: catapult PR
on:
pull_request:
branches:
- main
push:
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.97
- 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.97
- 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.97
- 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
+18
View File
@@ -0,0 +1,18 @@
[package]
name = "catapult"
version = "0.3.5"
rust-version = "1.97"
license = "MIT"
description = "Service to carry out text message scheduling"
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.5", features = ["v4", "serde"] }
schedtxt_models = { git = "ssh://git@git.kundeng.us/phoenix/schedtxt_models.git", tag = "v0.5.3" }
swoosh = { git = "ssh://git@git.kundeng.us/phoenix/swoosh.git", tag = "v0.5.4" }
[dev-dependencies]
+44
View File
@@ -0,0 +1,44 @@
FROM rust:1.97 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/*
RUN mkdir -p -m 0700 ~/.ssh && \
ssh-keyscan git.kundeng.us >> ~/.ssh/known_hosts
COPY Cargo.toml Cargo.lock ./
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 src ./src
# If you have other directories like `templates` or `static`, copy them too
COPY .env ./.env
RUN --mount=type=ssh \
cargo build --release --quiet
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 --from=builder /usr/src/app/target/release/catapult .
COPY --from=builder /usr/src/app/.env .
# Set the command to run your application
# Ensure this matches the binary name copied above
CMD ["./catapult"]
+22
View File
@@ -0,0 +1,22 @@
Copyright (c) 2026 Kun Deng.
Permission is hereby granted, free of charge, to any person
obtaining a copy of this software and associated documentation
files (the "Software"), to deal in the Software without
restriction, including without limitation the rights to use,
copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the
Software is furnished to do so, subject to the following
conditions:
The above copyright notice and this permission notice shall be
included in all copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND,
EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES
OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND
NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT
HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY,
WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING
FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR
OTHER DEALINGS IN THE SOFTWARE.
-22
View File
@@ -1,22 +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/internal/version.Version=$(VERSION)' \
-X 'git.kundeng.us/phoenix/catapult/internal/version.BuildTime=$(BUILD_TIME)' \
-X 'git.kundeng.us/phoenix/catapult/internal/version.Commit=$(COMMIT)' \
-X 'git.kundeng.us/phoenix/catapult/internal/version.GoVersion=$(GO_VERSION)'" \
-o catapult cmd/catapult/main.go
.PHONY: install
install:
go install -ldflags="\
-X 'git.kundeng.us/phoenix/catapult/internal/version.Version=$(VERSION)' \
-X 'git.kundeng.us/phoenix/catapult/internal/version.BuildTime=$(BUILD_TIME)' \
-X 'git.kundeng.us/phoenix/catapult/internal/version.Commit=$(COMMIT)' \
-X 'git.kundeng.us/phoenix/catapult/internal/version.GoVersion=$(GO_VERSION)'"
-o catapult cmd/catapult/main.go
+27 -6
View File
@@ -1,8 +1,29 @@
# Catapult
A service used to send text messages
# catapult
A service that manages scheduling text messages.
Building the service
```
make build
```
## Getting started
Docker isn't required, but it is the quickest way to get started.
Copy over the `.env.docker.sample` to `.env`. From that point, make edits to the
file to get it working.
### Messaging configuration
The service vendor used to send the messages is `twilio`. After signing up and being
provided credentials, update the `TWILIO_*` variables.
### Client-Server interaction
In order for `catapult` to function properly, credentials need to be provided via
`SERVICE_USERNAME` and `SERVICE_PASSPHRASE`. The credentials are created from making
an API call to register a service user in `textsender_auth`. So, you need to have credentials
set in the `.env` file and use those same credentials to create the account in the
API call.
### API sources
The `AUTH_URL` and `API_URL` variables need to be set for the base URL for the
`textsender_auth` and `textsender_api` respectively.
Since this is meant to run along with `textsender_auth` and `textsender_api`, it isn't
meant to be run by itself. So no need to build or run the container in isolation, but
that capability is present.
-31
View File
@@ -1,31 +0,0 @@
package main
import (
"fmt"
"os"
"os/signal"
"syscall"
"git.kundeng.us/phoenix/catapult/internal/config"
"git.kundeng.us/phoenix/catapult/internal/version"
)
func main() {
fmt.Println(config.App_Name)
versionFlag := config.CheckVersionFlag()
if *versionFlag {
fmt.Println(version.String())
return
}
service := config.NewService()
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
service.Start()
<-sigChan
service.Stop()
}
+11
View File
@@ -0,0 +1,11 @@
version: '3.8' # Use a recent version
services:
songparser:
build:
context: .
ssh: ["default"]
container_name: catapult
env_file:
- .env
restart: unless-stopped
-3
View File
@@ -1,3 +0,0 @@
module git.kundeng.us/phoenix/catapult
go 1.25.4
-102
View File
@@ -1,102 +0,0 @@
package config
import (
"context"
"flag"
"log"
"sync"
"time"
)
const App_Name = "catapult"
func CheckVersionFlag() *bool {
versionFlag := flag.Bool("version", false, "Print version information")
flag.Parse()
return versionFlag
}
type Service struct {
wg sync.WaitGroup
ctx context.Context
cancel context.CancelFunc
}
func NewService() *Service {
ctx, cancel := context.WithCancel(context.Background())
return &Service{
ctx: ctx,
cancel: cancel,
}
}
func (s *Service) Start() {
log.Println("Starting service...")
// Start multiple background workers
s.wg.Add(3)
go s.worker1()
go s.worker2()
go s.healthChecker()
}
func (s *Service) worker1() {
defer s.wg.Done()
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case <-s.ctx.Done():
log.Println("Worker 1 shutting down...")
return
case <-ticker.C:
// Do some work
log.Println("Worker 1: Processing...")
}
}
}
func (s *Service) worker2() {
defer s.wg.Done()
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for {
select {
case <-s.ctx.Done():
log.Println("Worker 2 shutting down...")
return
case <-ticker.C:
// Do some work
log.Println("Worker 2: Processing...")
}
}
}
func (s *Service) healthChecker() {
defer s.wg.Done()
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
for {
select {
case <-s.ctx.Done():
log.Println("Health checker shutting down...")
return
case <-ticker.C:
log.Println("Health check: Service is healthy")
}
}
}
func (s *Service) Stop() {
log.Println("Shutting down service...")
s.cancel()
s.wg.Wait()
log.Println("Service stopped gracefully")
}
-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,
)
}
+60
View File
@@ -0,0 +1,60 @@
/// 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: schedtxt_models::config::auxiliary::TwilioConfig,
}
pub fn load_app() -> Result<App, std::io::Error> {
use schedtxt_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 {
let auth_url: String = auth_url_env.value;
let api_url: String = api_url_env.value;
let service_username: String = service_username_env.value;
let service_passphrase: String = service_passphrase_env.value;
match schedtxt_models::config::auxiliary::load_config() {
Ok(twilio_config) => Ok(App {
api_url,
auth_url,
service_username,
service_passphrase,
twilio_config,
}),
Err(err) => Err(err),
}
} else {
match reason {
Some(reason) => Err(std::io::Error::other(reason)),
None => Err(std::io::Error::other("No reason found")),
}
}
}
+3
View File
@@ -0,0 +1,3 @@
pub mod app;
pub mod service;
pub mod version;
+58
View File
@@ -0,0 +1,58 @@
#[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) => std::sync::Arc::new(app),
Err(err) => {
eprintln!("Error: {err:?}");
std::process::exit(-1);
}
};
println!("App: {app:?}");
let auth = catapult::service::auth::Auth {
app: std::sync::Arc::clone(&app),
};
let token = match auth.get_token().await {
Ok(token) => std::sync::Arc::new(token),
Err(err) => {
eprintln!("Error getting token");
eprintln!("Error: {err:?}");
std::process::exit(-1);
}
};
let mut svc = catapult::service::core::Service {
app: std::sync::Arc::clone(&app),
token: std::sync::Arc::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:?}");
let refreshed = std::sync::Arc::new(refresh);
svc.token = refreshed;
}
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: std::sync::Arc<crate::app::App>,
}
impl Auth {
pub async fn get_token(&self) -> Result<schedtxt_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: &schedtxt_models::token::Login,
) -> Result<schedtxt_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<schedtxt_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: schedtxt_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: &schedtxt_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: std::sync::Arc<crate::app::App>,
pub token: std::sync::Arc<schedtxt_models::token::LoginResult>,
}
impl Contact {
pub async fn get(
&self,
contact_id: &uuid::Uuid,
) -> Result<Vec<schedtxt_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<schedtxt_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<schedtxt_models::contact::Contact> = Vec::new();
for event in lrs.iter() {
let ll: schedtxt_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: std::sync::Arc<crate::app::App>,
pub token: std::sync::Arc<schedtxt_models::token::Login>,
}
impl Service {
pub async fn do_the_work(&mut self) {
println!("Checking queue");
use schedtxt_models::message::scheduling::ScheduledMessage;
let queue = super::queue::Queue {
app: std::sync::Arc::clone(&self.app),
token: std::sync::Arc::clone(&self.token),
};
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: std::sync::Arc::clone(&self.app),
token: std::sync::Arc::clone(&self.token),
};
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: std::sync::Arc::clone(&self.app),
token: std::sync::Arc::clone(&self.token),
};
let messages = super::message::Message {
app: std::sync::Arc::clone(&self.app),
token: std::sync::Arc::clone(&self.token),
};
let sch_msg = super::scheduler::ScheduledMessage {
app: std::sync::Arc::clone(&self.app),
token: std::sync::Arc::clone(&self.token),
};
let mer_req = super::mer::MessageEventResponse {
app: std::sync::Arc::clone(&self.app),
token: std::sync::Arc::clone(&self.token),
};
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 = schedtxt_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: &schedtxt_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: std::sync::Arc<crate::app::App>,
pub token: std::sync::Arc<schedtxt_models::token::LoginResult>,
}
impl Event {
pub async fn get(
&self,
scheduled_message_id: &uuid::Uuid,
) -> Result<Vec<schedtxt_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<schedtxt_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<
schedtxt_models::message::scheduling::ScheduledMessageEvent,
> = Vec::new();
for event in lrs.iter() {
let ll: schedtxt_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: std::sync::Arc<crate::app::App>,
pub token: std::sync::Arc<schedtxt_models::token::LoginResult>,
}
impl MessageEventResponse {
pub async fn record_message_event_response(
&self,
mer: &schedtxt_models::message::event::MessageEventResponse,
) -> Result<schedtxt_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": schedtxt_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<schedtxt_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: schedtxt_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: std::sync::Arc<crate::app::App>,
pub token: std::sync::Arc<schedtxt_models::token::LoginResult>,
}
impl Message {
pub async fn get(
&self,
message_id: &uuid::Uuid,
) -> Result<Vec<schedtxt_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<schedtxt_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<schedtxt_models::message::Message> = Vec::new();
for event in lrs.iter() {
let ll: schedtxt_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: std::sync::Arc<crate::app::App>,
pub token: std::sync::Arc<schedtxt_models::token::LoginResult>,
}
impl Queue {
pub async fn get_queue(
&self,
) -> Result<schedtxt_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<schedtxt_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: schedtxt_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: std::sync::Arc<crate::app::App>,
pub token: std::sync::Arc<schedtxt_models::token::LoginResult>,
}
impl ScheduledMessage {
pub async fn get(
&self,
scheduled_message_id: &uuid::Uuid,
) -> Result<schedtxt_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<schedtxt_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: schedtxt_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
}