Compare commits

2 Commits
Author SHA1 Message Date
Jason Ross 2423ff1bb0 initial mvp 2026-03-22 14:48:56 -05:00
Jason Ross f0485b7248 stable initial container state 2026-03-22 13:22:11 -05:00
9 changed files with 4620 additions and 37 deletions
+2 -1
View File
@@ -1 +1,2 @@
.env
.env
target/
+4039
View File
File diff suppressed because it is too large Load Diff
+18
View File
@@ -4,3 +4,21 @@ version = "0.1.0"
edition = "2021"
[dependencies]
anyhow = "1.0.102"
argon2 = "0.5"
aws-sdk-s3 = "1"
axum = { version = "0.7", features = ["ws", "macros"] }
axum-extra = { version = "0.9", features = ["typed-header"] }
chrono = { version = "0.4", features = ["serde"] }
futures = "0.3"
jsonwebtoken = "9"
redis = { version = "0.27", features = ["tokio-comp"] }
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
sqlx = { version = "0.8", features = ["postgres", "runtime-tokio-rustls", "uuid", "chrono", "json"] }
tokio = { version = "1", features = ["full"] }
tower = "0.5"
tower-http = { version = "0.6", features = ["cors", "trace", "compression-full", "request-id", "util"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
uuid = { version = "1", features = ["v4", "serde"] }
+192
View File
@@ -0,0 +1,192 @@
use crate::error::AppError;
use crate::state::SharedState;
use argon2::{
password_hash::{rand_core::OsRng, PasswordHash, PasswordHasher, PasswordVerifier, SaltString},
Argon2,
};
use axum::{
async_trait,
extract::{FromRef, FromRequestParts, State},
http::request::Parts,
Json,
};
use axum_extra::{
headers::{authorization::Bearer, Authorization},
TypedHeader,
};
use jsonwebtoken::{decode, encode, DecodingKey, EncodingKey, Header, Validation};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use sqlx::Row;
use chrono::{Utc, Duration};
#[derive(Debug, Serialize, Deserialize)]
pub struct Claims {
pub sub: Uuid,
pub exp: usize,
}
#[async_trait]
impl<S> FromRequestParts<S> for Claims
where
S: Send + Sync,
SharedState: axum::extract::FromRef<S>,
{
type Rejection = AppError;
async fn from_request_parts(parts: &mut Parts, state: &S) -> Result<Self, Self::Rejection> {
let TypedHeader(Authorization(bearer)) =
TypedHeader::<Authorization<Bearer>>::from_request_parts(parts, state)
.await
.map_err(|_| AppError::Auth("Missing or invalid authorization header".to_string()))?;
let app_state = SharedState::from_ref(state);
let token_data = decode::<Claims>(
bearer.token(),
&DecodingKey::from_secret(app_state.jwt_secret.as_bytes()),
&Validation::default(),
)
.map_err(|_| AppError::Auth("Invalid or expired token".to_string()))?;
Ok(token_data.claims)
}
}
#[derive(Deserialize)]
pub struct RegisterRequest {
pub username: String,
pub password: String,
pub display_name: String,
}
#[derive(Serialize)]
pub struct AuthResponse {
pub token: String,
}
pub async fn register(
State(state): State<SharedState>,
Json(payload): Json<RegisterRequest>,
) -> Result<Json<AuthResponse>, AppError> {
if payload.username.is_empty() || payload.password.is_empty() || payload.display_name.is_empty() {
return Err(AppError::BadRequest("Fields cannot be empty".to_string()));
}
let salt = SaltString::generate(&mut OsRng);
let argon2 = Argon2::default();
let password_hash = argon2
.hash_password(payload.password.as_bytes(), &salt)?
.to_string();
let user_id: Uuid = sqlx::query_scalar(
r#"
INSERT INTO users (username, display_name, password_hash)
VALUES ($1, $2, $3)
RETURNING id
"#,
)
.bind(&payload.username)
.bind(&payload.display_name)
.bind(&password_hash)
.fetch_one(&state.db)
.await
.map_err(|e| {
if let sqlx::Error::Database(db_err) = &e {
if db_err.constraint() == Some("users_username_key") {
return AppError::BadRequest("Username already exists".to_string());
}
}
AppError::from(e)
})?;
let token = generate_token(user_id, &state.jwt_secret)?;
Ok(Json(AuthResponse { token }))
}
#[derive(Deserialize)]
pub struct LoginRequest {
pub username: String,
pub password: String,
}
pub async fn login(
State(state): State<SharedState>,
Json(payload): Json<LoginRequest>,
) -> Result<Json<AuthResponse>, AppError> {
let row = sqlx::query(
r#"
SELECT id, password_hash
FROM users
WHERE username = $1
"#,
)
.bind(&payload.username)
.fetch_optional(&state.db)
.await?;
if let Some(row) = row {
let id: Uuid = row.try_get("id")?;
let hash: String = row.try_get("password_hash")?;
let parsed_hash = PasswordHash::new(&hash)?;
let argon2 = Argon2::default();
if argon2
.verify_password(payload.password.as_bytes(), &parsed_hash)
.is_ok()
{
let token = generate_token(id, &state.jwt_secret)?;
return Ok(Json(AuthResponse { token }));
}
}
Err(AppError::Auth("Invalid username or password".to_string()))
}
#[derive(Serialize)]
pub struct UserProfile {
pub id: Uuid,
pub username: String,
pub display_name: String,
pub avatar_url: Option<String>,
}
pub async fn me(
claims: Claims,
State(state): State<SharedState>,
) -> Result<Json<UserProfile>, AppError> {
let profile = sqlx::query_as!(
UserProfile,
r#"
SELECT id, username, display_name, avatar_url
FROM users
WHERE id = $1
"#,
claims.sub
)
.fetch_optional(&state.db)
.await?
.ok_or_else(|| AppError::Auth("User not found".to_string()))?;
Ok(Json(profile))
}
fn generate_token(user_id: Uuid, secret: &str) -> Result<String, AppError> {
let expiration = Utc::now()
.checked_add_signed(Duration::days(7))
.expect("valid timestamp")
.timestamp() as usize;
let claims = Claims {
sub: user_id,
exp: expiration,
};
encode(
&Header::default(),
&claims,
&EncodingKey::from_secret(secret.as_bytes()),
)
.map_err(|e| AppError::Internal(e.into()))
}
+54
View File
@@ -0,0 +1,54 @@
use axum::{
http::StatusCode,
response::{IntoResponse, Response},
Json,
};
use serde_json::json;
#[derive(Debug)]
pub enum AppError {
Internal(anyhow::Error),
Auth(String),
BadRequest(String),
}
impl IntoResponse for AppError {
fn into_response(self) -> Response {
let (status, error_message) = match self {
AppError::Internal(err) => {
tracing::error!("Internal server error: {:?}", err);
(StatusCode::INTERNAL_SERVER_ERROR, "Internal server error".to_string())
}
AppError::Auth(msg) => {
(StatusCode::UNAUTHORIZED, msg.clone())
}
AppError::BadRequest(msg) => {
(StatusCode::BAD_REQUEST, msg.clone())
}
};
let body = Json(json!({
"error": error_message,
}));
(status, body).into_response()
}
}
impl From<sqlx::Error> for AppError {
fn from(inner: sqlx::Error) -> Self {
AppError::Internal(anyhow::anyhow!(inner))
}
}
impl From<anyhow::Error> for AppError {
fn from(inner: anyhow::Error) -> Self {
AppError::Internal(inner)
}
}
impl From<argon2::password_hash::Error> for AppError {
fn from(inner: argon2::password_hash::Error) -> Self {
AppError::Internal(anyhow::anyhow!(inner))
}
}
+97 -7
View File
@@ -1,9 +1,99 @@
use std::time::Duration;
use std::thread;
mod auth;
mod error;
mod state;
mod ws;
fn main() {
println!("Backend stub started...");
loop {
thread::sleep(Duration::from_secs(60));
}
use axum::{
routing::{get, post},
Router,
};
use state::{AppState, SharedState};
use std::{collections::HashMap, sync::Arc};
use tokio::net::TcpListener;
use tokio::sync::RwLock;
use tower_http::{
compression::CompressionLayer,
cors::CorsLayer,
request_id::{MakeRequestUuid, PropagateRequestIdLayer, SetRequestIdLayer},
trace::{DefaultMakeSpan, DefaultOnResponse, TraceLayer},
};
use tracing::Level;
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::registry()
.with(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "backend=debug,tower_http=debug,axum=debug".into()),
)
.with(tracing_subscriber::fmt::layer())
.init();
let database_url = std::env::var("DATABASE_URL")
.unwrap_or_else(|_| "postgres://postgres:postgres@localhost:5432/postcard".to_string());
let db = sqlx::postgres::PgPoolOptions::new()
.max_connections(5)
.connect(&database_url)
.await?;
let jwt_secret = std::env::var("JWT_SECRET").unwrap_or_else(|_| "super-secret-key".to_string());
let state: SharedState = Arc::new(AppState {
db,
jwt_secret,
connections: Arc::new(RwLock::new(HashMap::new())),
});
let app = Router::new()
.route("/api/auth/register", post(auth::register))
.route("/api/auth/login", post(auth::login))
.route("/api/auth/me", get(auth::me))
.route("/ws", get(ws::ws_handler))
.with_state(state)
.layer(
TraceLayer::new_for_http()
.make_span_with(DefaultMakeSpan::new().include_headers(true))
.on_response(DefaultOnResponse::new().level(Level::INFO)),
)
.layer(PropagateRequestIdLayer::x_request_id())
.layer(SetRequestIdLayer::x_request_id(MakeRequestUuid))
.layer(CorsLayer::permissive())
.layer(CompressionLayer::new());
let listener = TcpListener::bind("0.0.0.0:3000").await?;
tracing::info!("Listening on {}", listener.local_addr()?);
let server = axum::serve(listener, app)
.with_graceful_shutdown(shutdown_signal());
server.await?;
Ok(())
}
async fn shutdown_signal() {
let ctrl_c = async {
tokio::signal::ctrl_c()
.await
.expect("failed to install Ctrl+C handler");
};
#[cfg(unix)]
let terminate = async {
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.expect("failed to install signal handler")
.recv()
.await;
};
#[cfg(not(unix))]
let terminate = std::future::pending::<()>();
tokio::select! {
_ = ctrl_c => {},
_ = terminate => {},
}
tracing::info!("Received termination signal, shutting down...");
}
+17
View File
@@ -0,0 +1,17 @@
use sqlx::PgPool;
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::{mpsc, RwLock};
use uuid::Uuid;
use crate::ws::OutboundFrame;
pub type UserId = Uuid;
pub struct AppState {
pub db: PgPool,
pub jwt_secret: String,
pub connections: Arc<RwLock<HashMap<UserId, Vec<mpsc::Sender<OutboundFrame>>>>>,
}
pub type SharedState = Arc<AppState>;
+181
View File
@@ -0,0 +1,181 @@
use crate::auth::Claims;
use crate::state::{SharedState, UserId};
use axum::{
extract::{
ws::{Message, WebSocket, WebSocketUpgrade},
State,
},
response::IntoResponse,
};
use futures::{sink::SinkExt, stream::StreamExt};
use jsonwebtoken::{decode, DecodingKey, Validation};
use serde::{Deserialize, Serialize};
use std::time::Duration;
use tokio::{sync::mpsc, time::timeout};
use chrono::Utc;
#[derive(Deserialize, Debug)]
#[serde(tag = "type")]
pub enum InboundFrame {
Authenticate { token: String },
#[serde(other)]
Unknown,
}
#[derive(Serialize, Debug, Clone)]
#[serde(tag = "type")]
pub enum OutboundFrame {
PresenceUpdate {
user_id: UserId,
status: String,
#[serde(skip_serializing_if = "Option::is_none")]
last_seen: Option<String>,
},
Error {
message: String,
},
}
pub async fn ws_handler(
ws: WebSocketUpgrade,
State(state): State<SharedState>,
) -> impl IntoResponse {
ws.on_upgrade(move |socket| handle_socket(socket, state))
}
async fn handle_socket(mut socket: WebSocket, state: SharedState) {
let auth_msg = match timeout(Duration::from_secs(5), socket.next()).await {
Ok(Some(Ok(msg))) => msg,
_ => return,
};
let token = if let Message::Text(text) = auth_msg {
if let Ok(InboundFrame::Authenticate { token }) = serde_json::from_str(&text) {
token
} else {
let err_json = serde_json::to_string(&OutboundFrame::Error { message: "Expected Authenticate frame".into() }).unwrap();
let _ = socket.send(Message::Text(err_json)).await;
return;
}
} else {
return;
};
let user_id = match decode::<Claims>(
&token,
&DecodingKey::from_secret(state.jwt_secret.as_bytes()),
&Validation::default(),
) {
Ok(token_data) => token_data.claims.sub,
Err(_) => {
let err_json = serde_json::to_string(&OutboundFrame::Error { message: "Invalid token".into() }).unwrap();
let _ = socket.send(Message::Text(err_json)).await;
return;
}
};
let (tx, mut rx) = mpsc::channel::<OutboundFrame>(32);
{
let mut connections = state.connections.write().await;
connections.entry(user_id).or_default().push(tx);
}
let presence_msg = OutboundFrame::PresenceUpdate {
user_id,
status: "online".to_string(),
last_seen: None,
};
broadcast_all(&state, presence_msg).await;
let (mut sender, mut receiver) = socket.split();
let mut ping_interval = tokio::time::interval(Duration::from_secs(30));
ping_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let pong_timeout = tokio::time::sleep(Duration::MAX);
tokio::pin!(pong_timeout);
let mut awaiting_pong = false;
loop {
tokio::select! {
_ = ping_interval.tick() => {
if sender.send(Message::Ping(vec![])).await.is_err() {
break;
}
awaiting_pong = true;
pong_timeout.as_mut().reset(tokio::time::Instant::now() + Duration::from_secs(10));
}
_ = &mut pong_timeout, if awaiting_pong => {
break;
}
msg = rx.recv() => {
if let Some(frame) = msg {
if let Ok(json) = serde_json::to_string(&frame) {
if sender.send(Message::Text(json)).await.is_err() {
break;
}
}
} else {
break;
}
}
msg = receiver.next() => {
match msg {
Some(Ok(Message::Text(text))) => {
if let Ok(frame) = serde_json::from_str::<InboundFrame>(&text) {
match frame {
InboundFrame::Unknown | InboundFrame::Authenticate { .. } => {
let err_json = serde_json::to_string(&OutboundFrame::Error { message: "Unknown or unexpected frame type".into() }).unwrap();
let _ = sender.send(Message::Text(err_json)).await;
}
}
} else {
let err_json = serde_json::to_string(&OutboundFrame::Error { message: "Invalid JSON".into() }).unwrap();
let _ = sender.send(Message::Text(err_json)).await;
}
}
Some(Ok(Message::Pong(_))) => {
awaiting_pong = false;
pong_timeout.as_mut().reset(tokio::time::Instant::now() + Duration::from_secs(86400));
}
Some(Ok(Message::Close(_))) => {
break;
}
Some(Err(_)) | None => break,
_ => {}
}
}
}
}
let mut connections = state.connections.write().await;
let mut is_empty = false;
if let Some(user_conns) = connections.get_mut(&user_id) {
drop(rx);
user_conns.retain(|tx| !tx.is_closed());
if user_conns.is_empty() {
is_empty = true;
}
}
if is_empty {
connections.remove(&user_id);
drop(connections);
let offline_msg = OutboundFrame::PresenceUpdate {
user_id,
status: "offline".to_string(),
last_seen: Some(Utc::now().to_rfc3339()),
};
broadcast_all(&state, offline_msg).await;
}
}
async fn broadcast_all(state: &SharedState, frame: OutboundFrame) {
let connections = state.connections.read().await;
for conns in connections.values() {
for tx in conns {
let _ = tx.send(frame.clone()).await;
}
}
}
+20 -29
View File
@@ -4,14 +4,14 @@ services:
postgres:
image: postgres:18
environment:
POSTGRES_USER: ${POSTGRES_USER:-admin}
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:-Lord5673}
POSTGRES_DB: ${POSTGRES_DB:-postcard-im-app}
POSTGRES_USER: ${POSTGRES_USER}
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD}
POSTGRES_DB: ${POSTGRES_DB}
volumes:
- postgres_data:/var/lib/postgresql/data
- postgres_data:/var/lib/postgresql
- ./migrations:/docker-entrypoint-initdb.d:ro
healthcheck:
test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-admin} -d ${POSTGRES_DB:-postcard-im-app}"]
test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-admin} -d ${POSTGRES_DB}"]
interval: 5s
timeout: 5s
retries: 5
@@ -31,8 +31,8 @@ services:
minio:
image: minio/minio:latest
environment:
MINIO_ROOT_USER: ${MINIO_ROOT_USER:-minioadmin}
MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD:-ChuckNorris52!}
MINIO_ROOT_USER: ${MINIO_ROOT_USER}
MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD}
volumes:
- minio_data:/data
command: server /data --console-address ":9001"
@@ -59,16 +59,16 @@ services:
backend:
build:
context: ./backend
context: ./api
dockerfile: Containerfile
environment:
- DATABASE_URL=postgres://${POSTGRES_USER:-admin}:${POSTGRES_PASSWORD:-Lord5673}@postgres:5432/${POSTGRES_DB:-postcard-im-app}
- REDIS_URL=${REDIS_URL:-redis://redis:6379}
- MINIO_ENDPOINT=${MINIO_ENDPOINT:-http://minio:9000}
- MINIO_ROOT_USER=${MINIO_ROOT_USER:-minioadmin}
- MINIO_ROOT_PASSWORD=${MINIO_ROOT_PASSWORD:-ChuckNorris52!}
- MINIO_BUCKET=${MINIO_BUCKET:-media}
- JWT_SECRET=${JWT_SECRET:-changeme}
- DATABASE_URL=postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB}
- REDIS_URL=${REDIS_URL}
- MINIO_ENDPOINT=${MINIO_ENDPOINT}
- MINIO_ROOT_USER=${MINIO_ROOT_USER}
- MINIO_ROOT_PASSWORD=${MINIO_ROOT_PASSWORD}
- MINIO_BUCKET=${MINIO_BUCKET}
- JWT_SECRET=${JWT_SECRET}
ports:
- "8080:8080"
depends_on:
@@ -78,23 +78,11 @@ services:
condition: service_healthy
minio:
condition: service_healthy
web:
build:
context: ./web
dockerfile: Containerfile
environment:
- BACKEND_URL=${BACKEND_URL:-http://backend:8080}
ports:
- "3000:3000"
depends_on:
- backend
caddy:
image: caddy:latest
ports:
- "80:80"
- "443:443"
- "8000:80"
- "8443:443"
volumes:
- ./Caddyfile:/etc/caddy/Caddyfile
- caddy_data:/data
@@ -108,3 +96,6 @@ volumes:
minio_data:
caddy_data:
caddy_config:
secrets:
env_file:
file: .env