added get telemetry
This commit is contained in:
@@ -2,4 +2,6 @@ FROM rust:1-bookworm
|
||||
RUN apt update && apt upgrade -y && apt install fish iputils-ping -y
|
||||
RUN rustup component add clippy rustfmt
|
||||
|
||||
RUN useradd -ms /bin/fish vscode
|
||||
RUN useradd -ms /bin/fish vscode
|
||||
USER vscode
|
||||
RUN cargo install sqlx-cli
|
||||
@@ -7,7 +7,7 @@ edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
actix-web = "4.4.0"
|
||||
chrono = "0.4.31"
|
||||
chrono = { version = "0.4.31", features = ["serde"] }
|
||||
env_logger = "0.10.0"
|
||||
log = "0.4.20"
|
||||
serde = "1.0.188"
|
||||
|
||||
@@ -1,10 +1,10 @@
|
||||
-- Add migration script here
|
||||
CREATE TABLE Telemetry (
|
||||
timestamp TIMESTAMP NOT NULL,
|
||||
Software_Version INT,
|
||||
Software_Version INT NOT NULL,
|
||||
Voltage FLOAT,
|
||||
Temperature FLOAT,
|
||||
uptime INT,
|
||||
device_id CHAR(32),
|
||||
uptime INT NOT NULL,
|
||||
device_id CHAR(32) NOT NULL,
|
||||
FOREIGN KEY (device_id) REFERENCES Devices(ID)
|
||||
);
|
||||
@@ -1,10 +1,10 @@
|
||||
use actix_web::{cookie::time::error, web};
|
||||
use log::{error, info};
|
||||
use sqlx::{pool, postgres::PgPoolOptions, query, PgPool, Pool, Postgres, migrate};
|
||||
use sqlx::{pool, postgres::PgPoolOptions, query, PgPool, Pool, Postgres, migrate, query_as};
|
||||
use thiserror::Error;
|
||||
use chrono::Utc;
|
||||
|
||||
use crate::schemas::TelemetryMessage;
|
||||
use crate::schemas::{TelemetryMessage, TelemetryMessageFromDevice};
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Database {
|
||||
@@ -78,13 +78,25 @@ impl Database {
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn add_telemetry(&self, msg: &web::Json<TelemetryMessage>, device_id: &str) -> Result<(), DatabaseError> {
|
||||
pub async fn add_telemetry(&self, msg: &web::Json<TelemetryMessageFromDevice>, device_id: &str) -> Result<(), DatabaseError> {
|
||||
info!("Adding telemetry message to DB");
|
||||
let current_timestamp = Utc::now().naive_utc();
|
||||
query!("INSERT INTO Telemetry (timestamp, software_version, voltage, temperature, uptime, device_id)
|
||||
VALUES ($1, $2, $3, $4, $5, $6);",
|
||||
current_timestamp, msg.version, msg.voltage, msg.temperature, msg.uptime, device_id
|
||||
current_timestamp, msg.software_version, msg.voltage, msg.temperature, msg.uptime, device_id
|
||||
).execute(&self.conn_pool).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn get_telemetry_for_id(&self, device_id: &str) -> Result<Vec<TelemetryMessage>, DatabaseError> {
|
||||
info!("Getting telemetry messages for {} from DB", &device_id);
|
||||
let messages = query_as!(TelemetryMessage,
|
||||
"SELECT timestamp, software_version, voltage, temperature, uptime
|
||||
FROM Telemetry
|
||||
WHERE device_id = $1 ORDER BY timestamp DESC;", &device_id)
|
||||
.fetch_all(&self.conn_pool).await?;
|
||||
|
||||
|
||||
Ok(messages)
|
||||
}
|
||||
}
|
||||
|
||||
36
src/main.rs
36
src/main.rs
@@ -1,7 +1,7 @@
|
||||
use actix_web::{post, web, App, HttpServer, Responder, http::StatusCode, HttpResponse};
|
||||
use actix_web::{post, web, App, HttpServer, Responder, http::StatusCode, HttpResponse, get};
|
||||
use database::Database;
|
||||
use log::{info, error, debug};
|
||||
use crate::schemas::TelemetryMessage;
|
||||
use crate::schemas::TelemetryMessageFromDevice;
|
||||
|
||||
mod database;
|
||||
mod schemas;
|
||||
@@ -14,7 +14,7 @@ struct AppState {
|
||||
async fn receive_telemetry(
|
||||
device_id: web::Path<String>,
|
||||
data: web::Data<AppState>,
|
||||
telemetry_message: web::Json<TelemetryMessage>
|
||||
telemetry_message: web::Json<TelemetryMessageFromDevice>
|
||||
) -> impl Responder {
|
||||
info!("POST - telementry - Processing device id {}", device_id);
|
||||
match data.db.create_device_if_not_exists(&device_id).await{
|
||||
@@ -24,10 +24,33 @@ async fn receive_telemetry(
|
||||
return HttpResponse::InternalServerError();
|
||||
}
|
||||
};
|
||||
debug!("{:?}", telemetry_message);
|
||||
data.db.add_telemetry(&telemetry_message, &device_id).await;
|
||||
|
||||
HttpResponse::Ok()
|
||||
match data.db.add_telemetry(&telemetry_message, &device_id).await{
|
||||
Ok(_) => HttpResponse::Created(),
|
||||
Err(e) => {
|
||||
error!("adding Telemetry message to DB failed \n{}", e);
|
||||
HttpResponse::InternalServerError()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[get("/telemetry/{device_id}")]
|
||||
async fn get_telemetry(
|
||||
device_id: web::Path<String>,
|
||||
data: web::Data<AppState>
|
||||
) -> impl Responder {
|
||||
info!("GET - telementry - Processing device id {}", device_id);
|
||||
let messages = match data.db.get_telemetry_for_id(&device_id).await{
|
||||
Ok(msgs) => msgs,
|
||||
Err(e) => {
|
||||
error!("Getting Telemetry Messages from DB failed \n{}", e);
|
||||
return HttpResponse::InternalServerError().finish()
|
||||
}
|
||||
};
|
||||
|
||||
|
||||
HttpResponse::Ok().json(messages)
|
||||
|
||||
}
|
||||
|
||||
#[actix_web::main]
|
||||
@@ -43,6 +66,7 @@ async fn main() -> std::io::Result<()> {
|
||||
App::new()
|
||||
.app_data(web::Data::new(AppState { db: db.clone() }))
|
||||
.service(receive_telemetry)
|
||||
.service(get_telemetry)
|
||||
})
|
||||
.bind(("127.0.0.1", 8080))?
|
||||
.run()
|
||||
|
||||
@@ -1,9 +1,19 @@
|
||||
use serde::Deserialize;
|
||||
use chrono::NaiveDateTime;
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
#[derive(Deserialize, Debug, Serialize)]
|
||||
pub struct TelemetryMessage {
|
||||
pub uptime: i32,
|
||||
pub voltage: Option<f64>,
|
||||
pub temperature: Option<f64>,
|
||||
pub version: i32
|
||||
pub software_version: i32,
|
||||
pub timestamp: NaiveDateTime
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug, Serialize)]
|
||||
pub struct TelemetryMessageFromDevice {
|
||||
pub uptime: i32,
|
||||
pub voltage: Option<f64>,
|
||||
pub temperature: Option<f64>,
|
||||
pub software_version: i32,
|
||||
}
|
||||
Reference in New Issue
Block a user