Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ec690cd23e | ||
|
|
2880e4142f | ||
|
|
0c2d638adf | ||
|
|
0691ce0adb | ||
|
|
40d92ab7d3 | ||
|
|
7ec1f58f6a | ||
|
|
66c7fc261f | ||
|
|
a4fefd5bc6 | ||
|
|
c809815549 |
@@ -0,0 +1,8 @@
|
||||
target/
|
||||
|
||||
.git/
|
||||
|
||||
.gitea/
|
||||
.github/
|
||||
|
||||
*.DS_Store
|
||||
@@ -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
|
||||
@@ -0,0 +1,70 @@
|
||||
name: catapult PR
|
||||
|
||||
on:
|
||||
pull_request:
|
||||
branches:
|
||||
- main
|
||||
|
||||
concurrency:
|
||||
group: ${{ github.workflow }}-${{ github.ref }}
|
||||
cancel-in-progress: true
|
||||
|
||||
jobs:
|
||||
check:
|
||||
name: Check
|
||||
runs-on: ubuntu-24.04
|
||||
steps:
|
||||
- uses: actions/checkout@v5
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
toolchain: 1.96
|
||||
- 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@v5
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
toolchain: 1.96
|
||||
- run: rustup component add rustfmt
|
||||
- uses: Swatinem/rust-cache@v2
|
||||
- run: |
|
||||
mkdir -p ~/.ssh
|
||||
echo "${{ secrets.MYREPO_TOKEN }}" > ~/.ssh/textsender-models_deploy_key
|
||||
chmod 600 ~/.ssh/textsender-models_deploy_key
|
||||
ssh-keyscan ${{ secrets.MY_HOST }} >> ~/.ssh/known_hosts
|
||||
|
||||
eval $(ssh-agent -s)
|
||||
ssh-add -v ~/.ssh/textsender-models_deploy_key
|
||||
cargo fmt --all -- --check
|
||||
|
||||
clippy:
|
||||
name: Clippy
|
||||
runs-on: ubuntu-24.04
|
||||
steps:
|
||||
- uses: actions/checkout@v5
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
toolchain: 1.96
|
||||
- run: rustup component add 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
|
||||
@@ -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 ./...
|
||||
+2
-3
@@ -1,5 +1,4 @@
|
||||
catapult
|
||||
/vendor
|
||||
|
||||
/target
|
||||
.env
|
||||
.env.docker
|
||||
.env.local
|
||||
|
||||
Generated
+3245
File diff suppressed because it is too large
Load Diff
+18
@@ -0,0 +1,18 @@
|
||||
[package]
|
||||
name = "catapult"
|
||||
version = "0.3.0"
|
||||
rust-version = "1.96"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
tokio = { version = "1.52.3", features = ["full"] }
|
||||
futures = { version = "0.3.32" }
|
||||
reqwest = { version = "0.13.4", features = ["json", "stream", "multipart"] }
|
||||
serde = { version = "1.0.228", features = ["derive"] }
|
||||
serde_json = { version = "1.0.150" }
|
||||
time = { version = "0.3.49", features = ["formatting", "macros", "parsing", "serde"] }
|
||||
uuid = { version = "1.23.3", features = ["v4", "serde"] }
|
||||
textsender_models = { git = "ssh://git@git.kundeng.us/phoenix/textsender_models.git", tag = "v0.4.10" }
|
||||
swoosh = { git = "ssh://git@git.kundeng.us/phoenix/swoosh.git", tag = "v0.4.2.0" }
|
||||
|
||||
[dev-dependencies]
|
||||
+61
@@ -0,0 +1,61 @@
|
||||
# Stage 1: Build the application
|
||||
FROM rust:1.96 as builder
|
||||
|
||||
# Set the working directory inside the container
|
||||
WORKDIR /usr/src/app
|
||||
|
||||
# Install build dependencies if needed (e.g., git for cloning)
|
||||
RUN apt-get update && apt-get install -y --no-install-recommends \
|
||||
pkg-config libssl3 \
|
||||
ca-certificates \
|
||||
openssh-client git \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
# << --- ADD HOST KEY HERE --- >>
|
||||
# Replace 'yourgithost.com' with the actual hostname (e.g., github.com)
|
||||
RUN mkdir -p -m 0700 ~/.ssh && \
|
||||
ssh-keyscan git.kundeng.us >> ~/.ssh/known_hosts
|
||||
|
||||
# Copy Cargo manifests
|
||||
COPY Cargo.toml Cargo.lock ./
|
||||
|
||||
# Build *only* dependencies to leverage Docker cache
|
||||
# This dummy build caches dependencies as a separate layer
|
||||
RUN --mount=type=ssh mkdir src && \
|
||||
echo "fn main() {println!(\"if you see this, the build broke\")}" > src/main.rs && \
|
||||
cargo build --release --quiet && \
|
||||
rm -rf src target/release/deps/catapult*
|
||||
|
||||
# Copy the actual source code
|
||||
COPY src ./src
|
||||
# If you have other directories like `templates` or `static`, copy them too
|
||||
COPY .env ./.env
|
||||
|
||||
# << --- SSH MOUNT ADDED HERE --- >>
|
||||
# Build *only* dependencies to leverage Docker cache
|
||||
# This dummy build caches dependencies as a separate layer
|
||||
# Mount the SSH agent socket for this command
|
||||
RUN --mount=type=ssh \
|
||||
cargo build --release --quiet
|
||||
|
||||
# Stage 2: Create the final, smaller runtime image
|
||||
# Use a minimal base image like debian-slim or even distroless for security/size
|
||||
FROM debian:trixie-slim
|
||||
|
||||
# Install runtime dependencies if needed (e.g., SSL certificates)
|
||||
RUN apt-get update && apt-get install -y ca-certificates libssl-dev libssl3 && rm -rf /var/lib/apt/lists/*
|
||||
|
||||
# Set the working directory
|
||||
WORKDIR /usr/local/bin
|
||||
|
||||
# Copy the compiled binary from the builder stage
|
||||
# Replace 'catapult' with the actual name of your binary (usually the crate name)
|
||||
COPY --from=builder /usr/src/app/target/release/catapult .
|
||||
|
||||
# Copy other necessary files like .env (if used for runtime config) or static assets
|
||||
# It's generally better to configure via environment variables in Docker though
|
||||
COPY --from=builder /usr/src/app/.env .
|
||||
|
||||
# Set the command to run your application
|
||||
# Ensure this matches the binary name copied above
|
||||
CMD ["./catapult"]
|
||||
@@ -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
|
||||
@@ -1,8 +0,0 @@
|
||||
# Catapult
|
||||
A service used to send text messages
|
||||
|
||||
|
||||
Building the service
|
||||
```
|
||||
make build
|
||||
```
|
||||
@@ -1,38 +0,0 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"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() {
|
||||
fmt.Println(app.App_Name)
|
||||
|
||||
versionFlag := config.CheckVersionFlag()
|
||||
if *versionFlag {
|
||||
fmt.Println(version.String())
|
||||
return
|
||||
}
|
||||
|
||||
myApp, err := app.Load()
|
||||
if err != nil {
|
||||
fmt.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()
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
version: '3.8' # Use a recent version
|
||||
|
||||
services:
|
||||
songparser:
|
||||
build:
|
||||
context: .
|
||||
ssh: ["default"] # Uses host's SSH agent
|
||||
container_name: catapult
|
||||
env_file:
|
||||
- .env
|
||||
restart: unless-stopped # Optional: Restart policy
|
||||
@@ -1,17 +0,0 @@
|
||||
module git.kundeng.us/phoenix/catapult
|
||||
|
||||
go 1.25.4
|
||||
|
||||
require (
|
||||
git.kundeng.us/phoenix/swoosh v0.0.5
|
||||
git.kundeng.us/phoenix/textsender-models v0.0.10
|
||||
github.com/google/uuid v1.6.0
|
||||
github.com/joho/godotenv v1.5.1
|
||||
)
|
||||
|
||||
require (
|
||||
github.com/golang-jwt/jwt/v5 v5.3.0 // indirect
|
||||
github.com/golang/mock v1.6.0 // indirect
|
||||
github.com/pkg/errors v0.9.1 // indirect
|
||||
github.com/twilio/twilio-go v1.28.7 // indirect
|
||||
)
|
||||
@@ -1,60 +0,0 @@
|
||||
git.kundeng.us/phoenix/swoosh v0.0.5 h1:++kc6xRlt01ToMeIOSzYtGnYjug6x0p0euZwJ+xHDs4=
|
||||
git.kundeng.us/phoenix/swoosh v0.0.5/go.mod h1:BTn6gDaHWTP/OpvDjOFZJwN9kJ1xd9Rl6yFJGy8xV1k=
|
||||
git.kundeng.us/phoenix/textsender-models v0.0.10 h1:0AMe/zm4Xte2BcuqqhRP7+42m+4To+rUnHkR2a9eXlI=
|
||||
git.kundeng.us/phoenix/textsender-models v0.0.10/go.mod h1:9iPDQJg1Tc6WMNoW5+f8YKmnosMwlWHJ++hmxNLDEe0=
|
||||
github.com/beevik/etree v1.1.0/go.mod h1:r8Aw8JqVegEf0w2fDnATrX9VpkMcyFeM0FhwO62wh+A=
|
||||
github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E=
|
||||
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
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.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk=
|
||||
github.com/golang-jwt/jwt/v5 v5.3.0 h1:pv4AsKCKKZuqlgs5sUmn4x8UlGa0kEVt/puTpKx9vvo=
|
||||
github.com/golang-jwt/jwt/v5 v5.3.0/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/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
|
||||
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
|
||||
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
|
||||
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/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno=
|
||||
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/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
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.28.7 h1:WzzQDR/rqmNkVs1TwtcHFPYGTdSdPnx/eAZf5UIXzr4=
|
||||
github.com/twilio/twilio-go v1.28.7/go.mod h1:FpgNWMoD8CFnmukpKq9RNpUSGXC0BwnbeKZj2YHlIkw=
|
||||
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/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c h1:dUUwHk2QECo/6vqA44rthZ8ie2QXMNeKRTHCNY2nXvo=
|
||||
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
@@ -1,81 +0,0 @@
|
||||
package app
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"path"
|
||||
|
||||
"git.kundeng.us/phoenix/textsender-models/tx0/config"
|
||||
"github.com/joho/godotenv"
|
||||
)
|
||||
|
||||
const App_Name = "catapult"
|
||||
|
||||
type App struct {
|
||||
ApiUrl string
|
||||
AuthUrl string
|
||||
ServiceUsername string
|
||||
ServicePassphrase string
|
||||
TwilioConfig *config.TwiloConfig
|
||||
}
|
||||
|
||||
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() (*config.TwiloConfig, 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 := config.TwiloConfig{}
|
||||
cfg.AccountSID = authSid
|
||||
cfg.AuthToken = authToken
|
||||
cfg.ServiceSID = serviceSid
|
||||
cfg.Number = phoneNumber
|
||||
return &cfg, nil
|
||||
}
|
||||
}
|
||||
@@ -1,12 +0,0 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"flag"
|
||||
)
|
||||
|
||||
func CheckVersionFlag() *bool {
|
||||
versionFlag := flag.Bool("version", false, "Print version information")
|
||||
flag.Parse()
|
||||
|
||||
return versionFlag
|
||||
}
|
||||
@@ -1,55 +0,0 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
|
||||
"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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,178 +0,0 @@
|
||||
package core
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"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...")
|
||||
if s.token == nil {
|
||||
catapultAuth := service.Auth{Application: s.App}
|
||||
if token, err := catapultAuth.GetToken(); err != nil {
|
||||
fmt.Println("Error:", err)
|
||||
} else {
|
||||
fmt.Println("Access token:", token.AccessToken)
|
||||
s.token = token
|
||||
}
|
||||
}
|
||||
|
||||
log.Println("Twilio config account sid:", s.App.TwilioConfig.AccountSID)
|
||||
log.Println("Twilio config auth token:", s.App.TwilioConfig.AuthToken)
|
||||
log.Println("Twilio config service sid:", s.App.TwilioConfig.ServiceSID)
|
||||
log.Println("Twilio config number:", s.App.TwilioConfig.Number)
|
||||
|
||||
queue := service.Queue{Application: s.App, Token: s.token}
|
||||
if item, exists, err := queue.GetQueue(); err != nil {
|
||||
fmt.Println("Error:", err)
|
||||
} else {
|
||||
if *exists {
|
||||
fmt.Println("Scheduled message Id:", item.Id)
|
||||
fmt.Println("Created:", item.Created)
|
||||
fmt.Println("Scheduled:", item.Scheduled)
|
||||
fmt.Println("Status:", item.Status)
|
||||
log.Println("User Id:", item.UserId)
|
||||
if scheduledMessageValid(item, now) {
|
||||
fmt.Println("Scheduled Message can be sent")
|
||||
scheduler := service.Scheduler{Application: s.App, Token: s.token}
|
||||
if events, err := scheduler.GetEvents(item.Id); err != nil {
|
||||
fmt.Println("Error getting event:", err)
|
||||
} else {
|
||||
for _, event := range events {
|
||||
if msg, err := scheduler.GetMessage(event.MessageId); err != nil {
|
||||
fmt.Println("Error getting message:", err)
|
||||
} else {
|
||||
if c, err := scheduler.GetContact(event.RecipientId); err != nil {
|
||||
fmt.Println("Error getting contact:", err)
|
||||
} else {
|
||||
fmt.Println("Message Id:", msg.Id)
|
||||
fmt.Println("Contact Id:", c.Id)
|
||||
if res, err := scheduler.ScheduleMessage(*item, *event, *c, *msg); err != nil {
|
||||
log.Println("Failure with scheduling the message:", err)
|
||||
} else {
|
||||
if res != nil {
|
||||
bytes, err := json.Marshal(res)
|
||||
if err != nil {
|
||||
log.Println("Error parsing result:", err)
|
||||
} else {
|
||||
jsonValue := string(bytes)
|
||||
log.Println("Result:", jsonValue)
|
||||
}
|
||||
} else {
|
||||
log.Println("Result should not be empty")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
fmt.Println("Invalid scheduled message")
|
||||
}
|
||||
} else {
|
||||
fmt.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
|
||||
}
|
||||
}
|
||||
@@ -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,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,59 +0,0 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/url"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"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 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
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,58 +0,0 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"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"
|
||||
"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, error) {
|
||||
msgSender := send.MessageSender{Config: s.Application.TwilioConfig}
|
||||
if res, err := msgSender.Send(msg, c, &schMsg.Scheduled); err != nil {
|
||||
return nil, err
|
||||
} else {
|
||||
fmt.Println(res.AccountSid)
|
||||
return res, nil
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
)
|
||||
}
|
||||
+56
@@ -0,0 +1,56 @@
|
||||
/// Seconds to sleep after finising a task
|
||||
pub const SECONDS_TO_SLEEP: u64 = 5;
|
||||
|
||||
pub const APP_NAME: &str = "catapult";
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct App {
|
||||
pub api_url: String,
|
||||
pub auth_url: String,
|
||||
pub service_username: String,
|
||||
pub service_passphrase: String,
|
||||
pub twilio_config: textsender_models::config::auxiliary::TwilioConfig,
|
||||
}
|
||||
|
||||
// TODO: Change this to non-async at some point.
|
||||
// Will require changes in other libraries
|
||||
pub async fn load_app() -> Result<App, std::io::Error> {
|
||||
use textsender_models::envy;
|
||||
let auth_url_env = envy::environment::get_env("AUTH_URL").await;
|
||||
let api_url_env = envy::environment::get_env("API_URL").await;
|
||||
let service_username_env = envy::environment::get_env("SERVICE_USERNAME").await;
|
||||
let service_passphrase_env = envy::environment::get_env("SERVICE_PASSPHRASE").await;
|
||||
|
||||
let envs = vec![
|
||||
api_url_env,
|
||||
auth_url_env,
|
||||
service_username_env,
|
||||
service_passphrase_env,
|
||||
];
|
||||
|
||||
let check_envs = |vs: &Vec<envy::EnvVar>| -> (bool, Option<String>) {
|
||||
for env in vs {
|
||||
if env.value.is_empty() {
|
||||
return (false, Some(format!("Key {} is not provided", env.key)));
|
||||
}
|
||||
}
|
||||
(true, None)
|
||||
};
|
||||
|
||||
let (envs_valid, reason) = check_envs(&envs);
|
||||
|
||||
if envs_valid {
|
||||
match textsender_models::config::auxiliary::load_config().await {
|
||||
Ok(twilio_config) => Ok(App {
|
||||
api_url: envs[0].value.clone(),
|
||||
auth_url: envs[1].value.clone(),
|
||||
service_username: envs[2].value.clone(),
|
||||
service_passphrase: envs[3].value.clone(),
|
||||
twilio_config,
|
||||
}),
|
||||
Err(err) => Err(err),
|
||||
}
|
||||
} else {
|
||||
Err(std::io::Error::other(reason.unwrap()))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
pub mod app;
|
||||
pub mod service;
|
||||
pub mod version;
|
||||
+55
@@ -0,0 +1,55 @@
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let args = std::env::args().collect();
|
||||
if catapult::version::contains_version(&args) {
|
||||
catapult::version::print_version();
|
||||
std::process::exit(-1);
|
||||
}
|
||||
|
||||
let app = match catapult::app::load_app().await {
|
||||
Ok(app) => app,
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
std::process::exit(-1);
|
||||
}
|
||||
};
|
||||
|
||||
println!("App: {app:?}");
|
||||
|
||||
let auth = catapult::service::auth::Auth { app: app.clone() };
|
||||
let token = match auth.get_token().await {
|
||||
Ok(token) => token,
|
||||
Err(err) => {
|
||||
eprintln!("Error getting token");
|
||||
eprintln!("Error: {err:?}");
|
||||
std::process::exit(-1);
|
||||
}
|
||||
};
|
||||
|
||||
let mut svc = catapult::service::core::Service {
|
||||
app: app.clone(),
|
||||
token,
|
||||
};
|
||||
|
||||
loop {
|
||||
if catapult::service::auth::has_token_expired(&svc.token).await {
|
||||
match auth.get_refresh_token(&svc.token).await {
|
||||
Ok(refresh) => {
|
||||
println!("Refresh token: {refresh:?}");
|
||||
svc.token = refresh;
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
std::process::exit(-1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
svc.do_the_work().await;
|
||||
|
||||
tokio::time::sleep(tokio::time::Duration::from_secs(
|
||||
catapult::app::SECONDS_TO_SLEEP,
|
||||
))
|
||||
.await;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,116 @@
|
||||
pub struct Auth {
|
||||
pub app: crate::app::App,
|
||||
}
|
||||
|
||||
impl Auth {
|
||||
pub async fn get_token(&self) -> Result<textsender_models::token::Login, std::io::Error> {
|
||||
let client = reqwest::Client::new();
|
||||
let endpoint = String::from("api/v1/service/login");
|
||||
let api_url = format!("{}/{endpoint}", self.app.auth_url);
|
||||
|
||||
println!("Url: {api_url:?}");
|
||||
|
||||
let payload = serde_json::json!({
|
||||
"username": &self.app.service_username,
|
||||
"passphrase": &self.app.service_passphrase,
|
||||
});
|
||||
|
||||
println!("Payload: {payload:?}");
|
||||
|
||||
match client.post(api_url).json(&payload).send().await {
|
||||
Ok(response) => match response.status() {
|
||||
reqwest::StatusCode::OK => {
|
||||
println!("Response: {response:?}");
|
||||
match response.text().await {
|
||||
Ok(val) => match parse_token_response(&val).await {
|
||||
Ok(login_result) => Ok(login_result),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err)),
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error getting token")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn get_refresh_token(
|
||||
&self,
|
||||
login_result: &textsender_models::token::Login,
|
||||
) -> Result<textsender_models::token::Login, std::io::Error> {
|
||||
let client = reqwest::Client::new();
|
||||
let endpoint = String::from("api/v1/token/refresh");
|
||||
let api_url = format!("{}/{endpoint}", self.app.auth_url);
|
||||
|
||||
println!("Url: {api_url:?}");
|
||||
|
||||
let payload = serde_json::json!({
|
||||
"access_token": login_result.access_token
|
||||
});
|
||||
|
||||
match client.post(api_url).json(&payload).send().await {
|
||||
Ok(response) => match response.status() {
|
||||
reqwest::StatusCode::OK => {
|
||||
println!("Response: {response:?}");
|
||||
match response.text().await {
|
||||
Ok(val) => match parse_token_response(&val).await {
|
||||
Ok(login_result) => Ok(login_result),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err)),
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error getting token")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn parse_token_response(
|
||||
response: &str,
|
||||
) -> Result<textsender_models::token::LoginResult, std::io::Error> {
|
||||
match serde_json::from_str::<serde_json::Value>(response) {
|
||||
Ok(j) => {
|
||||
println!("Good");
|
||||
let message = &j["message"];
|
||||
let data = &j["data"];
|
||||
|
||||
println!("Message: {message:?}");
|
||||
println!("Data: {data:?}");
|
||||
match j.get("data") {
|
||||
Some(serde_json::Value::Array(lrs)) => {
|
||||
if lrs.is_empty() {
|
||||
Err(std::io::Error::other("Error response is empty"))
|
||||
} else {
|
||||
let lr = lrs[0].clone();
|
||||
let ll: textsender_models::token::LoginResult =
|
||||
match serde_json::from_value(lr) {
|
||||
Ok(lr) => lr,
|
||||
Err(err) => {
|
||||
return Err(std::io::Error::other(err.to_string()));
|
||||
}
|
||||
};
|
||||
Ok(ll)
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error parsing response")),
|
||||
}
|
||||
}
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn has_token_expired(token: &textsender_models::token::LoginResult) -> bool {
|
||||
let now = time::OffsetDateTime::now_utc();
|
||||
let expire = match time::OffsetDateTime::from_unix_timestamp(token.expires_in) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
return false;
|
||||
}
|
||||
};
|
||||
|
||||
now > expire
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
pub struct Contact {
|
||||
pub app: crate::app::App,
|
||||
pub token: textsender_models::token::LoginResult,
|
||||
}
|
||||
|
||||
impl Contact {
|
||||
pub async fn get(
|
||||
&self,
|
||||
contact_id: &uuid::Uuid,
|
||||
) -> Result<Vec<textsender_models::contact::Contact>, std::io::Error> {
|
||||
let client = reqwest::Client::new();
|
||||
let endpoint = format!("api/v1/contact?id={contact_id}");
|
||||
let api_url = format!("{}/{endpoint}", self.app.api_url);
|
||||
|
||||
println!("Url: {api_url:?}");
|
||||
|
||||
let (key, header) = super::auth_header(&self.token.access_token);
|
||||
|
||||
match client.get(api_url).header(key, header).send().await {
|
||||
Ok(response) => match response.status() {
|
||||
reqwest::StatusCode::OK => match response.text().await {
|
||||
Ok(val) => match parse_response(&val).await {
|
||||
Ok(login_result) => Ok(login_result),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err)),
|
||||
},
|
||||
_ => Err(std::io::Error::other("Error getting token")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn parse_response(
|
||||
response: &str,
|
||||
) -> Result<Vec<textsender_models::contact::Contact>, std::io::Error> {
|
||||
match serde_json::from_str::<serde_json::Value>(response) {
|
||||
Ok(j) => match j.get("data") {
|
||||
Some(serde_json::Value::Array(lrs)) => {
|
||||
if lrs.is_empty() {
|
||||
Err(std::io::Error::other("Error response is empty"))
|
||||
} else {
|
||||
let mut events: Vec<textsender_models::contact::Contact> = Vec::new();
|
||||
|
||||
for event in lrs.iter() {
|
||||
let ll: textsender_models::contact::Contact =
|
||||
match serde_json::from_value(event.clone()) {
|
||||
Ok(lr) => lr,
|
||||
Err(err) => {
|
||||
return Err(std::io::Error::other(err.to_string()));
|
||||
}
|
||||
};
|
||||
events.push(ll);
|
||||
}
|
||||
|
||||
Ok(events)
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error parsing response")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,177 @@
|
||||
pub struct Service {
|
||||
pub app: crate::app::App,
|
||||
pub token: textsender_models::token::Login,
|
||||
}
|
||||
|
||||
impl Service {
|
||||
pub async fn do_the_work(&mut self) {
|
||||
println!("Do some work");
|
||||
println!("Checking queue");
|
||||
|
||||
use textsender_models::message::scheduling::ScheduledMessage;
|
||||
|
||||
let queue = super::queue::Queue {
|
||||
app: self.app.clone(),
|
||||
token: self.token.clone(),
|
||||
};
|
||||
let queue_item: Option<ScheduledMessage> = match queue.get_queue().await {
|
||||
Ok(item) => Some(item),
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
None
|
||||
}
|
||||
};
|
||||
|
||||
if queue_item.is_none() {
|
||||
println!("No queue item found");
|
||||
return;
|
||||
}
|
||||
let queue_item = queue_item.unwrap();
|
||||
|
||||
println!("Queue item Id: {:?}", queue_item.id);
|
||||
|
||||
if !self.is_schedulable(&queue_item) {
|
||||
println!("Not schedulable");
|
||||
return;
|
||||
}
|
||||
|
||||
println!("Can schedule");
|
||||
|
||||
let events_req = super::event::Event {
|
||||
app: self.app.clone(),
|
||||
token: self.token.clone(),
|
||||
};
|
||||
|
||||
let scheduled_message_id = queue_item.id;
|
||||
let events = match events_req.get(&scheduled_message_id).await {
|
||||
Ok(events) => Some(events),
|
||||
Err(err) => {
|
||||
eprintln!("Error: {err:?}");
|
||||
None
|
||||
}
|
||||
};
|
||||
|
||||
let contacts = super::contact::Contact {
|
||||
app: self.app.clone(),
|
||||
token: self.token.clone(),
|
||||
};
|
||||
let messages = super::message::Message {
|
||||
app: self.app.clone(),
|
||||
token: self.token.clone(),
|
||||
};
|
||||
let sch_msg = super::scheduler::ScheduledMessage {
|
||||
app: self.app.clone(),
|
||||
token: self.token.clone(),
|
||||
};
|
||||
let mer_req = super::mer::MessageEventResponse {
|
||||
app: self.app.clone(),
|
||||
token: self.token.clone(),
|
||||
};
|
||||
|
||||
let events = events.unwrap_or_default();
|
||||
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:?}");
|
||||
|
||||
match swoosh::twilio::api::send_message(message, contact, ¶m, twilio_config).await {
|
||||
Ok(resp) => {
|
||||
println!("Message scheduled");
|
||||
// MER response
|
||||
let result = swoosh::twilio::api::response_to_json(resp).await;
|
||||
let mer = textsender_models::message::event::MessageEventResponse {
|
||||
contact_id: contact.id.unwrap(),
|
||||
scheduled_message_event_id: Some(event.id),
|
||||
response: result,
|
||||
message_id: message.id.unwrap(),
|
||||
sent: scheduled_message.scheduled,
|
||||
user_id: scheduled_message.user_id,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
match mer_req.record_message_event_response(&mer).await {
|
||||
Ok(response) => {
|
||||
println!("Recording MER");
|
||||
println!("Response: {response:?}");
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Error recording mer");
|
||||
eprintln!("Error: {err:?}");
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
eprintln!("Message not sent");
|
||||
eprintln!("Error: {err:?}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
todo!(
|
||||
r#"
|
||||
Need to finish this
|
||||
|
||||
[X] - This function should fetch a token, check if the token is valid, and obtain a refresh token if
|
||||
it has expired.
|
||||
|
||||
Then do the real work.
|
||||
|
||||
1. Check the queue - Done
|
||||
2. If there is an item, get the data - Done
|
||||
3. Evaluate if it is schedulable - Done
|
||||
4. Get the events and iterate each one - Done
|
||||
5. Starting with the message - Done
|
||||
6. Get the contact - Done
|
||||
7. Send the message
|
||||
8. Record the message event response - Added, not used
|
||||
9. Add to scheduler
|
||||
"#
|
||||
);
|
||||
*/
|
||||
}
|
||||
|
||||
fn is_schedulable(
|
||||
&self,
|
||||
queue_item: &textsender_models::message::scheduling::ScheduledMessage,
|
||||
) -> bool {
|
||||
let now = time::OffsetDateTime::now_utc();
|
||||
|
||||
match queue_item.scheduled {
|
||||
Some(date) => date > now,
|
||||
None => false,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,70 @@
|
||||
pub struct Event {
|
||||
pub app: crate::app::App,
|
||||
pub token: textsender_models::token::LoginResult,
|
||||
}
|
||||
|
||||
impl Event {
|
||||
pub async fn get(
|
||||
&self,
|
||||
scheduled_message_id: &uuid::Uuid,
|
||||
) -> Result<Vec<textsender_models::message::scheduling::ScheduledMessageEvent>, std::io::Error>
|
||||
{
|
||||
let client = reqwest::Client::new();
|
||||
let endpoint =
|
||||
format!("api/v1/schedule/message/event?scheduled_message_id={scheduled_message_id}");
|
||||
let api_url = format!("{}/{endpoint}", self.app.api_url);
|
||||
|
||||
println!("Url: {api_url:?}");
|
||||
|
||||
let (key, header) = super::auth_header(&self.token.access_token);
|
||||
|
||||
match client.get(api_url).header(key, header).send().await {
|
||||
Ok(response) => match response.status() {
|
||||
reqwest::StatusCode::OK => match response.text().await {
|
||||
Ok(val) => match parse_response(&val).await {
|
||||
Ok(login_result) => Ok(login_result),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err)),
|
||||
},
|
||||
_ => Err(std::io::Error::other(
|
||||
"Error getting Scheduled Message Event",
|
||||
)),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn parse_response(
|
||||
response: &str,
|
||||
) -> Result<Vec<textsender_models::message::scheduling::ScheduledMessageEvent>, std::io::Error> {
|
||||
match serde_json::from_str::<serde_json::Value>(response) {
|
||||
Ok(j) => match j.get("data") {
|
||||
Some(serde_json::Value::Array(lrs)) => {
|
||||
if lrs.is_empty() {
|
||||
Err(std::io::Error::other("Error response is empty"))
|
||||
} else {
|
||||
let mut events: Vec<
|
||||
textsender_models::message::scheduling::ScheduledMessageEvent,
|
||||
> = Vec::new();
|
||||
|
||||
for event in lrs.iter() {
|
||||
let ll: textsender_models::message::scheduling::ScheduledMessageEvent =
|
||||
match serde_json::from_value(event.clone()) {
|
||||
Ok(lr) => lr,
|
||||
Err(err) => {
|
||||
return Err(std::io::Error::other(err.to_string()));
|
||||
}
|
||||
};
|
||||
events.push(ll);
|
||||
}
|
||||
|
||||
Ok(events)
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error parsing response")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
pub struct MessageEventResponse {
|
||||
pub app: crate::app::App,
|
||||
pub token: textsender_models::token::LoginResult,
|
||||
}
|
||||
|
||||
impl MessageEventResponse {
|
||||
pub async fn record_message_event_response(
|
||||
&self,
|
||||
mer: &textsender_models::message::event::MessageEventResponse,
|
||||
) -> Result<textsender_models::message::event::MessageEventResponse, std::io::Error> {
|
||||
let client = reqwest::Client::new();
|
||||
let endpoint = String::from("api/v1/schedule/message/event/response/record");
|
||||
let api_url = format!("{}/{endpoint}", self.app.api_url);
|
||||
|
||||
println!("Url: {api_url:?}");
|
||||
|
||||
use time::format_description::well_known::Iso8601;
|
||||
let sent = match mer.sent.unwrap().format(&Iso8601::DEFAULT) {
|
||||
Ok(converted) => converted,
|
||||
Err(err) => {
|
||||
return Err(std::io::Error::other(err.to_string()));
|
||||
}
|
||||
};
|
||||
|
||||
let payload = serde_json::json!({
|
||||
"scheduled_message_event_id": mer.scheduled_message_event_id.unwrap(),
|
||||
"response": mer.response,
|
||||
"user_id": mer.user_id,
|
||||
"sent": sent,
|
||||
"status": textsender_models::message::event::MESSAGE_EVENT_RESPONSE_STATUS_SCHEDULED,
|
||||
"contact_id": mer.contact_id,
|
||||
"message_id": mer.message_id
|
||||
});
|
||||
|
||||
println!("Payload: {payload:?}");
|
||||
|
||||
let (key, header) = super::auth_header(&self.token.access_token);
|
||||
|
||||
match client
|
||||
.post(api_url)
|
||||
.json(&payload)
|
||||
.header(key, header)
|
||||
.send()
|
||||
.await
|
||||
{
|
||||
Ok(response) => match response.status() {
|
||||
reqwest::StatusCode::CREATED | reqwest::StatusCode::OK => {
|
||||
println!("Response: {response:?}");
|
||||
match response.text().await {
|
||||
Ok(val) => match parse_response(&val).await {
|
||||
Ok(login_result) => Ok(login_result),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err)),
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other(
|
||||
"Error recording Message Event Response",
|
||||
)),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
*
|
||||
* {
|
||||
"scheduled_message_event_id": "{{scheduled_message_event_id}}",
|
||||
"response": {{text_response_raw}},
|
||||
"user_id": "{{user_id}}",
|
||||
"sent": "2026-06-20T17:10:00Z",
|
||||
"status": "{{mer_status_scheduled}}",
|
||||
"contact_id": "{{contact_id}}",
|
||||
"message_id": "{{message_id}}"
|
||||
}
|
||||
*
|
||||
*/
|
||||
}
|
||||
|
||||
async fn parse_response(
|
||||
response: &str,
|
||||
) -> Result<textsender_models::message::event::MessageEventResponse, std::io::Error> {
|
||||
match serde_json::from_str::<serde_json::Value>(response) {
|
||||
Ok(j) => {
|
||||
println!("Good");
|
||||
let message = &j["message"];
|
||||
let data = &j["data"];
|
||||
|
||||
println!("Message: {message:?}");
|
||||
println!("Data: {data:?}");
|
||||
match j.get("data") {
|
||||
Some(serde_json::Value::Array(lrs)) => {
|
||||
if lrs.is_empty() {
|
||||
Err(std::io::Error::other("Error response is empty"))
|
||||
} else {
|
||||
let lr = lrs[0].clone();
|
||||
let ll: textsender_models::message::event::MessageEventResponse =
|
||||
match serde_json::from_value(lr) {
|
||||
Ok(lr) => lr,
|
||||
Err(err) => {
|
||||
return Err(std::io::Error::other(err.to_string()));
|
||||
}
|
||||
};
|
||||
Ok(ll)
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error parsing response")),
|
||||
}
|
||||
}
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
pub struct Message {
|
||||
pub app: crate::app::App,
|
||||
pub token: textsender_models::token::LoginResult,
|
||||
}
|
||||
|
||||
impl Message {
|
||||
pub async fn get(
|
||||
&self,
|
||||
message_id: &uuid::Uuid,
|
||||
) -> Result<Vec<textsender_models::message::Message>, std::io::Error> {
|
||||
let client = reqwest::Client::new();
|
||||
let endpoint = format!("api/v1/message?id={message_id}");
|
||||
let api_url = format!("{}/{endpoint}", self.app.api_url);
|
||||
|
||||
println!("Url: {api_url:?}");
|
||||
|
||||
let (key, header) = super::auth_header(&self.token.access_token);
|
||||
|
||||
match client.get(api_url).header(key, header).send().await {
|
||||
Ok(response) => match response.status() {
|
||||
reqwest::StatusCode::OK => match response.text().await {
|
||||
Ok(val) => match parse_response(&val).await {
|
||||
Ok(login_result) => Ok(login_result),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err)),
|
||||
},
|
||||
_ => Err(std::io::Error::other("Error getting token")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn parse_response(
|
||||
response: &str,
|
||||
) -> Result<Vec<textsender_models::message::Message>, std::io::Error> {
|
||||
match serde_json::from_str::<serde_json::Value>(response) {
|
||||
Ok(j) => match j.get("data") {
|
||||
Some(serde_json::Value::Array(lrs)) => {
|
||||
if lrs.is_empty() {
|
||||
Err(std::io::Error::other("Error response is empty"))
|
||||
} else {
|
||||
let mut events: Vec<textsender_models::message::Message> = Vec::new();
|
||||
|
||||
for event in lrs.iter() {
|
||||
let ll: textsender_models::message::Message =
|
||||
match serde_json::from_value(event.clone()) {
|
||||
Ok(lr) => lr,
|
||||
Err(err) => {
|
||||
return Err(std::io::Error::other(err.to_string()));
|
||||
}
|
||||
};
|
||||
events.push(ll);
|
||||
}
|
||||
|
||||
Ok(events)
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error parsing response")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
pub struct Queue {
|
||||
pub app: crate::app::App,
|
||||
pub token: textsender_models::token::LoginResult,
|
||||
}
|
||||
|
||||
impl Queue {
|
||||
pub async fn get_queue(
|
||||
&self,
|
||||
) -> Result<textsender_models::message::scheduling::ScheduledMessage, std::io::Error> {
|
||||
let client = reqwest::Client::new();
|
||||
let endpoint = String::from("api/v1/schedule/message/fetch");
|
||||
let api_url = format!("{}/{endpoint}", self.app.api_url);
|
||||
|
||||
println!("Url: {api_url:?}");
|
||||
|
||||
let (key, header) = super::auth_header(&self.token.access_token);
|
||||
|
||||
match client.get(api_url).header(key, header).send().await {
|
||||
Ok(response) => match response.status() {
|
||||
reqwest::StatusCode::OK => {
|
||||
println!("Response: {response:?}");
|
||||
match response.text().await {
|
||||
Ok(val) => match parse_queue_response(&val).await {
|
||||
Ok(login_result) => Ok(login_result),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err)),
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error getting token")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn parse_queue_response(
|
||||
response: &str,
|
||||
) -> Result<textsender_models::message::scheduling::ScheduledMessage, std::io::Error> {
|
||||
match serde_json::from_str::<serde_json::Value>(response) {
|
||||
Ok(j) => {
|
||||
let message = &j["message"];
|
||||
let data = &j["data"];
|
||||
|
||||
println!("Message: {message:?}");
|
||||
println!("Data: {data:?}");
|
||||
match j.get("data") {
|
||||
Some(serde_json::Value::Array(lrs)) => {
|
||||
if lrs.is_empty() {
|
||||
Err(std::io::Error::other("Error response is empty"))
|
||||
} else {
|
||||
let lr = lrs[0].clone();
|
||||
let ll: textsender_models::message::scheduling::ScheduledMessage =
|
||||
match serde_json::from_value(lr) {
|
||||
Ok(lr) => lr,
|
||||
Err(err) => {
|
||||
return Err(std::io::Error::other(err.to_string()));
|
||||
}
|
||||
};
|
||||
Ok(ll)
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error parsing response")),
|
||||
}
|
||||
}
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
pub struct ScheduledMessage {
|
||||
pub app: crate::app::App,
|
||||
pub token: textsender_models::token::LoginResult,
|
||||
}
|
||||
|
||||
impl ScheduledMessage {
|
||||
pub async fn get(
|
||||
&self,
|
||||
scheduled_message_id: &uuid::Uuid,
|
||||
) -> Result<textsender_models::message::scheduling::ScheduledMessage, std::io::Error> {
|
||||
let client = reqwest::Client::new();
|
||||
let endpoint = format!("api/v1/schedule/message?id={scheduled_message_id}");
|
||||
let api_url = format!("{}/{endpoint}", self.app.api_url);
|
||||
|
||||
println!("Url: {api_url:?}");
|
||||
|
||||
let (key, header) = super::auth_header(&self.token.access_token);
|
||||
|
||||
match client.get(api_url).header(key, header).send().await {
|
||||
Ok(response) => match response.status() {
|
||||
reqwest::StatusCode::OK => match response.text().await {
|
||||
Ok(val) => match parse_response(&val).await {
|
||||
Ok(login_result) => Ok(login_result),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err)),
|
||||
},
|
||||
_ => Err(std::io::Error::other("Error getting Scheduled Message")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn parse_response(
|
||||
response: &str,
|
||||
) -> Result<textsender_models::message::scheduling::ScheduledMessage, std::io::Error> {
|
||||
match serde_json::from_str::<serde_json::Value>(response) {
|
||||
Ok(j) => match j.get("data") {
|
||||
Some(serde_json::Value::Array(lrs)) => {
|
||||
if lrs.is_empty() {
|
||||
Err(std::io::Error::other("Error response is empty"))
|
||||
} else {
|
||||
let ll: textsender_models::message::scheduling::ScheduledMessage =
|
||||
match serde_json::from_value(lrs[0].clone()) {
|
||||
Ok(lr) => lr,
|
||||
Err(err) => {
|
||||
return Err(std::io::Error::other(err.to_string()));
|
||||
}
|
||||
};
|
||||
|
||||
Ok(ll)
|
||||
}
|
||||
}
|
||||
_ => Err(std::io::Error::other("Error parsing response")),
|
||||
},
|
||||
Err(err) => Err(std::io::Error::other(err.to_string())),
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user