9 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
40 changed files with 3879 additions and 973 deletions
+5 -3
View File
@@ -1,6 +1,8 @@
vendor/
target/
.git/
.gitea/
.env
.gitea/
.github/
*.DS_Store
+1 -1
View File
@@ -4,5 +4,5 @@ SERVICE_USERNAME=suave
SERVICE_PASSPHRASE=9238urc9328nr329
TWILIO_AUTH_SID=9M438C93R943U4329MCU43C34U
TWILIO_SERVICE_SID=9M4J3X8439U398NUVT3342MC349C348T
TWILIO_AUTH_TOKEN="f4a1f2b0b79ea3735078c2d8ee9684e1"
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.26.1' # 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 ./...
+2 -4
View File
@@ -1,6 +1,4 @@
catapult
/vendor
/target
.env
.env.local
.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]
+26 -27
View File
@@ -1,45 +1,44 @@
FROM golang:1.26.1 AS builder
FROM rust:1.97 as builder
WORKDIR /app
# 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
openssh-client git \
&& rm -rf /var/lib/apt/lists/*
RUN mkdir -p -m 0700 ~/.ssh && \
ssh-keyscan git.kundeng.us >> ~/.ssh/known_hosts
# Configure Git to use SSH for GitHub
RUN git config --global url."ssh://git@git.kundeng.us".insteadOf "https://git.kundeng.us"
# Set up the Go environment for private modules
ENV GOPRIVATE=git.kundeng.us
# Copy go mod and sum files
COPY go.mod go.sum ./
COPY Cargo.toml Cargo.lock ./
RUN --mount=type=ssh mkdir src && \
go mod download
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 source code
COPY ./cmd ./cmd
COPY ./internal ./internal
COPY ./Makefile .
COPY ./.env .
COPY src ./src
# If you have other directories like `templates` or `static`, copy them too
COPY .env ./.env
# Build the application
RUN CGO_ENABLED=0 GOOS=linux make build
RUN --mount=type=ssh \
cargo build --release --quiet
# Runtime stage
FROM alpine:latest AS production
FROM debian:trixie-slim
RUN apk --no-cache add ca-certificates
# 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/*
WORKDIR /root/
# Set the working directory
WORKDIR /usr/local/bin
# Copy the pre-built binary file from the previous stage
COPY --from=builder /app/catapult .
COPY --from=builder /app/.env ./
COPY --from=builder /usr/src/app/target/release/catapult .
# Command to run the executable
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.
-39
View File
@@ -1,39 +0,0 @@
package main
import (
"fmt"
"log"
"os"
"os/signal"
"syscall"
"git.kundeng.us/phoenix/catapult/internal/app"
"git.kundeng.us/phoenix/catapult/internal/config"
"git.kundeng.us/phoenix/catapult/internal/service/core"
"git.kundeng.us/phoenix/catapult/internal/version"
)
func main() {
log.Println(app.App_Name)
versionFlag := config.CheckVersionFlag()
if *versionFlag {
fmt.Println(version.String())
return
}
myApp, err := app.Load()
if err != nil {
log.Println("Error loading app: %w", err)
return
}
service := core.NewService(myApp)
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
service.Start()
<-sigChan
service.Stop()
}
-11
View File
@@ -1,11 +0,0 @@
version: '3.8' # Use a recent version
services:
catapult_service:
build: # Tells docker-compose to build the Dockerfile in the current directory
context: .
ssh: ["default"] # Uses host's SSH agent
container_name: catapult # Optional: Give the container a specific name
env_file:
- .env
restart: unless-stopped # Optional: Restart policy
+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
-17
View File
@@ -1,17 +0,0 @@
module git.kundeng.us/phoenix/catapult
go 1.26.1
require (
git.kundeng.us/phoenix/swoosh v0.2.0
git.kundeng.us/phoenix/textsender-models v0.2.0
github.com/google/uuid v1.6.0
github.com/joho/godotenv v1.5.1
)
require (
github.com/golang-jwt/jwt/v5 v5.3.1 // indirect
github.com/golang/mock v1.6.0 // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/twilio/twilio-go v1.30.4 // indirect
)
-49
View File
@@ -1,49 +0,0 @@
git.kundeng.us/phoenix/swoosh v0.2.0 h1:aecJdtRhgQofaOfki0fm+YszUuwKyfKqF5YVb20yxHI=
git.kundeng.us/phoenix/swoosh v0.2.0/go.mod h1:NxWD9iUunLIJRgTvlZbZa6jyDHc0JPPRuTmsrxV/hMs=
git.kundeng.us/phoenix/textsender-models v0.2.0 h1:smz8Fs8VOs1Ya23txbOM0YPRidZIsM0yE9unHF0D/nQ=
git.kundeng.us/phoenix/textsender-models v0.2.0/go.mod h1:3CkqA/HFKPhpMYxkKn5uVbZEzEbG3sofLZE8pZ1BHO4=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY=
github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE=
github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc=
github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0=
github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4=
github.com/localtunnel/go-localtunnel v0.0.0-20170326223115-8a804488f275 h1:IZycmTpoUtQK3PD60UYBwjaCUHUP7cML494ao9/O8+Q=
github.com/localtunnel/go-localtunnel v0.0.0-20170326223115-8a804488f275/go.mod h1:zt6UU74K6Z6oMOYJbJzYpYucqdcQwSMPBEdSvGiaUMw=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/testify v1.7.0 h1:nwc3DEeHmmLAfoZucVR881uASk0Mfjw8xYJ99tb5CcY=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/twilio/twilio-go v1.30.4 h1:Whrz37IykDD9KJI2YX4LWaGxYWZoYB7va8fASGDrLng=
github.com/twilio/twilio-go v1.30.4/go.mod h1:QbitvbvtkV77Jn4BABAKVmxabYSjMyQG4tHey9gfPqg=
github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/mod v0.4.2/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20210405180319-a5a99cb37ef4/go.mod h1:p54w0d4576C0XHj96bSt6lcn1PtDYWL6XObtHCRCNQM=
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210510120138-977fb7262007/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.1.1/go.mod h1:o0xws9oXOQQZyjljx8fwUC0k7L1pTE6eaCbjGeHmOkk=
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
-81
View File
@@ -1,81 +0,0 @@
package app
import (
"fmt"
"os"
"path"
auxcfg "git.kundeng.us/phoenix/textsender-models/tx0/config/auxiliary"
"github.com/joho/godotenv"
)
const App_Name = "catapult"
type App struct {
ApiUrl string
AuthUrl string
ServiceUsername string
ServicePassphrase string
TwilioConfig *auxcfg.TwilioConfig
}
func Load() (*App, error) {
err := godotenv.Load()
if err != nil {
cwd, _ := os.Getwd()
envPath := path.Join(cwd, "../..", ".env")
if err = godotenv.Load(envPath); err != nil {
prevPath := path.Join(envPath, "../..", ".env")
if err = godotenv.Load(prevPath); err != nil {
return nil, fmt.Errorf("Error loading .env file: %w", err)
}
}
}
apiUrl := os.Getenv("API_URL")
authUrl := os.Getenv("AUTH_URL")
serviceUsername := os.Getenv("SERVICE_USERNAME")
servicePassphrase := os.Getenv("SERVICE_PASSPHRASE")
if len(apiUrl) == 0 {
return nil, fmt.Errorf("Api Url not provided")
} else if len(authUrl) == 0 {
return nil, fmt.Errorf("Auth url not provided")
} else if len(serviceUsername) == 0 {
return nil, fmt.Errorf("Service username not provided")
} else if len(servicePassphrase) == 0 {
return nil, fmt.Errorf("Service passphrase not provided")
} else {
if cfg, err := loadTwilioConfig(); err != nil {
return nil, err
} else {
return &App{
ApiUrl: apiUrl, AuthUrl: authUrl, ServiceUsername: serviceUsername, ServicePassphrase: servicePassphrase, TwilioConfig: cfg,
}, nil
}
}
}
func loadTwilioConfig() (*auxcfg.TwilioConfig, error) {
authSid := os.Getenv("TWILIO_AUTH_SID")
serviceSid := os.Getenv("TWILIO_SERVICE_SID")
authToken := os.Getenv("TWILIO_AUTH_TOKEN")
phoneNumber := os.Getenv("TWILIO_PHONE_NUMBER")
if len(authSid) == 0 {
return nil, fmt.Errorf("Twilio config auth sid not provided")
} else if len(serviceSid) == 0 {
return nil, fmt.Errorf("Twilio config service sid not provided")
} else if len(authToken) == 0 {
return nil, fmt.Errorf("Twilio config token not provided")
} else if len(phoneNumber) == 0 {
return nil, fmt.Errorf("Twilio config phone number not provided")
} else {
cfg := auxcfg.TwilioConfig{}
cfg.AccountSID = authSid
cfg.AuthToken = authToken
cfg.ServiceSID = serviceSid
cfg.Number = phoneNumber
return &cfg, nil
}
}
-12
View File
@@ -1,12 +0,0 @@
package config
import (
"flag"
)
func CheckVersionFlag() *bool {
versionFlag := flag.Bool("version", false, "Print version information")
flag.Parse()
return versionFlag
}
-98
View File
@@ -1,98 +0,0 @@
package service
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
"time"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type Auth struct {
Application *app.App
}
type service struct {
Username string `json:"username"`
Passphrase string `json:"passphrase"`
}
type tokenResponse struct {
Message string `json:"message"`
Data []*token.Login `json:"data"`
}
func (a *Auth) GetToken() (*token.Login, error) {
serv := service{
Username: a.Application.ServiceUsername,
Passphrase: a.Application.ServicePassphrase,
}
jsonData, err := json.Marshal(serv)
if err != nil {
return nil, err
}
resp, err := http.Post(
fmt.Sprintf("%s/api/v1/service/login", a.Application.AuthUrl),
"application/json",
bytes.NewBuffer(jsonData),
)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r tokenResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
}
return r.Data[0], nil
}
type refreshTokenRequest struct {
AccessToken string `json:"access_token"`
}
func (a *Auth) GetRefreshToken(tok *token.Login) (*token.Login, error) {
req := refreshTokenRequest{AccessToken: tok.AccessToken}
jsonData, err := json.Marshal(req)
if err != nil {
return nil, err
}
resp, err := http.Post(fmt.Sprintf("%s/api/v1/token/refresh", a.Application.AuthUrl), "application/json", bytes.NewBuffer(jsonData))
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r tokenResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data[0], nil
}
}
}
func (s *Auth) TokenExpired(tok *token.Login) bool {
now := time.Now()
expiredTime := time.Unix(tok.ExpiresIn, 0)
if now.After(expiredTime) {
return true
} else {
return false
}
}
-58
View File
@@ -1,58 +0,0 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"net/url"
"git.kundeng.us/phoenix/textsender-models/tx0/contact"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"github.com/google/uuid"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type Contact struct {
Application *app.App
Token *token.Login
}
type getContactResponse struct {
Message string `json:"message"`
Data []*contact.Contact `json:"data"`
}
func (c *Contact) GetContact(id uuid.UUID) (*contact.Contact, error) {
params := url.Values{}
params.Add("id", id.String())
pm := params.Encode()
fullUrl := fmt.Sprintf("%s/api/v1/contact?%s", c.Application.ApiUrl, pm)
fmt.Println("Url:", fullUrl)
req, err := http.NewRequest("GET", fullUrl, nil)
if err != nil {
return nil, err
}
token := c.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r getContactResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data[0], nil
}
}
}
-187
View File
@@ -1,187 +0,0 @@
package core
import (
"context"
"log"
"sync"
"time"
"git.kundeng.us/phoenix/textsender-models/tx0/message/scheduling"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"git.kundeng.us/phoenix/catapult/internal/app"
"git.kundeng.us/phoenix/catapult/internal/service"
)
type Service struct {
wg sync.WaitGroup
ctx context.Context
cancel context.CancelFunc
App *app.App
token *token.Login
}
func NewService(application *app.App) *Service {
ctx, cancel := context.WithCancel(context.Background())
return &Service{
ctx: ctx,
cancel: cancel,
App: application,
}
}
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
now := time.Now()
log.Println("Worker 1: Processing...")
catapultAuth := service.Auth{Application: s.App}
if s.token == nil {
log.Println("Token has not been fetched")
log.Println("Fetching token")
if token, err := catapultAuth.GetToken(); err != nil {
log.Println("Error:", err)
} else {
log.Println("Token fetched")
s.token = token
}
} else if catapultAuth.TokenExpired(s.token) {
// Get refresh token
log.Println("Token expired")
log.Println("Fetching refresh token")
refreshToken, err := catapultAuth.GetRefreshToken(s.token)
if err != nil {
log.Println("Error getting refresh token:", err)
} else {
log.Println("Refresh token fetched")
s.token = refreshToken
}
}
queue := service.Queue{Application: s.App, Token: s.token}
if item, exists, err := queue.GetQueue(); err != nil {
log.Println("Error:", err)
} else {
if *exists {
log.Println("Scheduled message Id:", item.Id)
log.Println("Created:", item.Created)
log.Println("Scheduled:", item.Scheduled)
log.Println("Status:", item.Status)
log.Println("User Id:", item.UserId)
if scheduledMessageValid(item, now) {
log.Println("Scheduled Message can be sent")
scheduler := service.Scheduler{Application: s.App, Token: s.token}
if events, err := scheduler.GetEvents(item.Id); err != nil {
log.Println("Error getting event:", err)
} else {
for _, event := range events {
if msg, err := scheduler.GetMessage(event.MessageId); err != nil {
log.Println("Error getting message:", err)
} else {
if c, err := scheduler.GetContact(event.ContactId); err != nil {
log.Println("Error getting contact:", err)
} else {
log.Println("Message Id:", msg.Id)
log.Println("Contact Id:", c.Id)
if res, resRaw, sent, err := scheduler.ScheduleMessage(*item, *event, *c, *msg); err != nil {
log.Println("Failure with scheduling the message:", err)
} else {
if res != nil {
if eventResponse, err := scheduler.RecordEventResponse(event, item.UserId, resRaw, sent); err != nil {
log.Println("Failure recording event response:", err)
} else {
if len(eventResponse) == 0 {
log.Println("No event responses")
} else {
eventResponse := eventResponse[0]
log.Println("Event response saved. Id:", eventResponse.Id)
}
}
} else {
log.Println("Result should not be empty")
}
}
}
}
}
}
} else {
log.Println("Invalid scheduled message")
}
} else {
log.Println("Empty queue")
}
}
}
}
}
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")
}
func scheduledMessageValid(schMsg *scheduling.ScheduledMessage, now time.Time) bool {
if schMsg.Scheduled.Before(now) {
return false
} else {
return true
}
}
-58
View File
@@ -1,58 +0,0 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"net/url"
"git.kundeng.us/phoenix/textsender-models/tx0/message"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"github.com/google/uuid"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type Message struct {
Application *app.App
Token *token.Login
}
type getMessageResponse struct {
Message string `json:"message"`
Data []*message.Message `json:"data"`
}
func (m *Message) GetMessage(id uuid.UUID) (*message.Message, error) {
params := url.Values{}
params.Add("id", id.String())
pm := params.Encode()
fullUrl := fmt.Sprintf("%s/api/v1/message?%s", m.Application.ApiUrl, pm)
fmt.Println("Url:", fullUrl)
req, err := http.NewRequest("GET", fullUrl, nil)
if err != nil {
return nil, err
}
token := m.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r getMessageResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data[0], nil
}
}
}
@@ -1,58 +0,0 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"strings"
"git.kundeng.us/phoenix/textsender-models/tx0/message/event"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type MessageEventResponse struct {
Application *app.App
Token *token.Login
}
type recordMessageEventResponse struct {
Message string `json:"message"`
Data []*event.MessageEventResponse `json:"data"`
}
func (m *MessageEventResponse) RecordEventResponse(mer *event.MessageEventResponse) ([]*event.MessageEventResponse, error) {
fullUrl := fmt.Sprintf("%s/api/v1/schedule/message/event/response/record", m.Application.ApiUrl)
jsonData, err := json.Marshal(mer)
if err != nil {
return nil, err
}
req, err := http.NewRequest("POST", fullUrl, strings.NewReader(string(jsonData)))
if err != nil {
return nil, err
}
token := m.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r recordMessageEventResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data, nil
}
}
}
-54
View File
@@ -1,54 +0,0 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"git.kundeng.us/phoenix/textsender-models/tx0/message/scheduling"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type Queue struct {
Application *app.App
Token *token.Login
}
type queueResponse struct {
Message string `json:"message"`
Data []*scheduling.ScheduledMessage `json:"data"`
}
func (q *Queue) GetQueue() (*scheduling.ScheduledMessage, *bool, error) {
req, err := http.NewRequest("GET", fmt.Sprintf("%s/api/v1/schedule/message/fetch", q.Application.ApiUrl), nil)
if err != nil {
return nil, nil, err
}
token := q.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, nil, err
}
defer resp.Body.Close()
var r queueResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, nil, err
} else {
var exists bool
if len(r.Data) == 0 {
exists = false
return nil, &exists, nil
} else {
exists = true
return r.Data[0], &exists, nil
}
}
}
@@ -1,58 +0,0 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"net/url"
"git.kundeng.us/phoenix/textsender-models/tx0/message/scheduling"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"github.com/google/uuid"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type ScheduledMessageEvent struct {
Application *app.App
Token *token.Login
}
type getScheduledMessageEventResponse struct {
Message string `json:"message"`
Data []*scheduling.ScheduledMessageEvent `json:"data"`
}
func (s *ScheduledMessageEvent) GetEvents(scheduledMessageId uuid.UUID) ([]*scheduling.ScheduledMessageEvent, error) {
params := url.Values{}
params.Add("scheduled_message_id", scheduledMessageId.String())
pm := params.Encode()
fullUrl := fmt.Sprintf("%s/api/v1/schedule/message/event?%s", s.Application.ApiUrl, pm)
req, err := http.NewRequest("GET", fullUrl, nil)
if err != nil {
return nil, err
}
token := s.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r getScheduledMessageEventResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data, nil
}
}
}
-69
View File
@@ -1,69 +0,0 @@
package service
import (
"time"
"git.kundeng.us/phoenix/swoosh/swoop/send"
"git.kundeng.us/phoenix/swoosh/swoop/types"
"git.kundeng.us/phoenix/textsender-models/tx0/contact"
"git.kundeng.us/phoenix/textsender-models/tx0/message"
evnt "git.kundeng.us/phoenix/textsender-models/tx0/message/event"
"git.kundeng.us/phoenix/textsender-models/tx0/message/scheduling"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"github.com/google/uuid"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type Scheduler struct {
Application *app.App
Token *token.Login
}
func (s *Scheduler) GetEvents(scheduledMessageId uuid.UUID) ([]*scheduling.ScheduledMessageEvent, error) {
sme := ScheduledMessageEvent{Application: s.Application, Token: s.Token}
if events, err := sme.GetEvents(scheduledMessageId); err != nil {
return nil, err
} else {
return events, nil
}
}
func (s *Scheduler) GetMessage(id uuid.UUID) (*message.Message, error) {
msg := Message{Application: s.Application, Token: s.Token}
if letter, err := msg.GetMessage(id); err != nil {
return nil, err
} else {
return letter, nil
}
}
func (s *Scheduler) GetContact(id uuid.UUID) (*contact.Contact, error) {
ctct := Contact{Application: s.Application, Token: s.Token}
if c, err := ctct.GetContact(id); err != nil {
return nil, err
} else {
return c, nil
}
}
func (s *Scheduler) ScheduleMessage(schMsg scheduling.ScheduledMessage, event scheduling.ScheduledMessageEvent, c contact.Contact, msg message.Message) (*types.TwilioResult, map[string]any, *time.Time, error) {
msgSender := send.MessageSender{Config: s.Application.TwilioConfig}
if res, resRaw, err := msgSender.Send(msg, c, &schMsg.Scheduled); err != nil {
return nil, nil, nil, err
} else {
sent := time.Now()
return res, resRaw, &sent, nil
}
}
func (s *Scheduler) RecordEventResponse(schMsgEvent *scheduling.ScheduledMessageEvent, userId uuid.UUID, bytes map[string]any, sent *time.Time) ([]*evnt.MessageEventResponse, error) {
mre := evnt.MessageEventResponse{ScheduledMessageEventId: schMsgEvent.Id, UserId: userId, Response: bytes, Sent: *sent, Status: evnt.Message_Event_Response_Status_Scheduled}
msgEventResponse := MessageEventResponse{Application: s.Application, Token: s.Token}
if responses, err := msgEventResponse.RecordEventResponse(&mre); err != nil {
return nil, err
} else {
return responses, nil
}
}
-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
}