Initial commit
This commit is contained in:
@@ -0,0 +1,57 @@
|
||||
package vassistent;
|
||||
|
||||
import vassistent.bootstrap.ApplicationContext;
|
||||
import vassistent.bootstrap.ApplicationInitializer;
|
||||
import vassistent.controller.AppWindowController;
|
||||
import vassistent.service.*;
|
||||
import vassistent.util.Logger;
|
||||
|
||||
import javax.swing.*;
|
||||
import java.io.File;
|
||||
|
||||
|
||||
public class App {
|
||||
|
||||
public static void main(String[] args) {
|
||||
|
||||
ApplicationContext context =
|
||||
ApplicationInitializer.initialize();
|
||||
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
|
||||
|
||||
Logger.info("SYSTEM", "Programm wird beendet");
|
||||
|
||||
deleteDatabase();
|
||||
|
||||
Logger.shutdown();
|
||||
|
||||
}));
|
||||
|
||||
SwingUtilities.invokeLater(() -> {
|
||||
new AppWindowController(context)
|
||||
.createAndShowWindow();
|
||||
});
|
||||
}
|
||||
|
||||
private static void deleteDatabase() {
|
||||
|
||||
try {
|
||||
|
||||
File dbFile = new File("data/health.db");
|
||||
|
||||
if (dbFile.exists()) {
|
||||
|
||||
boolean deleted = dbFile.delete();
|
||||
|
||||
if (deleted) {
|
||||
Logger.info("SYSTEM", "Datenbank gelöscht");
|
||||
} else {
|
||||
Logger.warn("SYSTEM", "Datenbank konnte nicht gelöscht werden");
|
||||
}
|
||||
}
|
||||
|
||||
} catch (Exception e) {
|
||||
Logger.error("SYSTEM", "DB Löschung fehlgeschlagen", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,81 @@
|
||||
package vassistent.bootstrap;
|
||||
|
||||
import vassistent.model.AppState;
|
||||
import vassistent.service.*;
|
||||
|
||||
public class ApplicationContext {
|
||||
|
||||
private AppState appState;
|
||||
|
||||
private DataPersistenceService persistenceService;
|
||||
private StatisticsService statisticsService;
|
||||
private EvaluationService evaluationService;
|
||||
private UnrealWebSocketService unrealService;
|
||||
private MqttClientService mqttService;
|
||||
private BinaryEventService binaryEventService;
|
||||
private ProcessManagerService processManagerService;
|
||||
|
||||
public void setAppState(AppState appState) {
|
||||
this.appState = appState;
|
||||
}
|
||||
|
||||
public void setPersistenceService(DataPersistenceService persistenceService) {
|
||||
this.persistenceService = persistenceService;
|
||||
}
|
||||
|
||||
public void setStatisticsService(StatisticsService statisticsService) {
|
||||
this.statisticsService = statisticsService;
|
||||
}
|
||||
|
||||
public void setEvaluationService(EvaluationService evaluationService) {
|
||||
this.evaluationService = evaluationService;
|
||||
}
|
||||
|
||||
public void setUnrealService(UnrealWebSocketService unrealService) {
|
||||
this.unrealService = unrealService;
|
||||
}
|
||||
|
||||
public void setMqttService(MqttClientService mqttService) {
|
||||
this.mqttService = mqttService;
|
||||
}
|
||||
|
||||
public void setBinaryEventService(BinaryEventService binaryEventService) {
|
||||
this.binaryEventService = binaryEventService;
|
||||
}
|
||||
|
||||
public AppState getAppState() {
|
||||
return appState;
|
||||
}
|
||||
|
||||
public DataPersistenceService getPersistenceService() {
|
||||
return persistenceService;
|
||||
}
|
||||
|
||||
public StatisticsService getStatisticsService() {
|
||||
return statisticsService;
|
||||
}
|
||||
|
||||
public EvaluationService getEvaluationService() {
|
||||
return evaluationService;
|
||||
}
|
||||
|
||||
public UnrealWebSocketService getUnrealService() {
|
||||
return unrealService;
|
||||
}
|
||||
|
||||
public MqttClientService getMqttService() {
|
||||
return mqttService;
|
||||
}
|
||||
|
||||
public BinaryEventService getBinaryEventService() {
|
||||
return binaryEventService;
|
||||
}
|
||||
|
||||
public ProcessManagerService getProcessManagerService() {
|
||||
return processManagerService;
|
||||
}
|
||||
|
||||
public void setProcessManagerService(ProcessManagerService processManagerService) {
|
||||
this.processManagerService = processManagerService;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
package vassistent.bootstrap;
|
||||
|
||||
import vassistent.controller.DashboardController;
|
||||
import vassistent.model.AppState;
|
||||
import vassistent.service.*;
|
||||
import vassistent.util.ConfigLoader;
|
||||
import vassistent.util.Logger;
|
||||
|
||||
import java.io.InputStream;
|
||||
import java.util.Properties;
|
||||
|
||||
public class ApplicationInitializer {
|
||||
|
||||
private static Properties config =
|
||||
ConfigLoader.loadProperties(
|
||||
"config/application.properties"
|
||||
);
|
||||
|
||||
public static ApplicationContext initialize() {
|
||||
|
||||
ApplicationContext context = new ApplicationContext();
|
||||
|
||||
// ===== Model State =====
|
||||
context.setAppState(new AppState());
|
||||
|
||||
// ===== Infrastructure Services =====
|
||||
context.setUnrealService(
|
||||
new UnrealWebSocketService("ws://localhost:8888/avatar")
|
||||
);
|
||||
|
||||
context.setMqttService(new MqttClientService(context.getAppState()));
|
||||
|
||||
context.setPersistenceService(new DataPersistenceService());
|
||||
|
||||
context.setStatisticsService(
|
||||
new StatisticsService(context.getPersistenceService())
|
||||
);
|
||||
|
||||
context.setEvaluationService(
|
||||
new EvaluationService(
|
||||
context.getStatisticsService(),
|
||||
context.getAppState(),
|
||||
context.getUnrealService()
|
||||
)
|
||||
);
|
||||
|
||||
context.setBinaryEventService(
|
||||
new BinaryEventService(
|
||||
context.getPersistenceService(),
|
||||
context.getEvaluationService()
|
||||
)
|
||||
);
|
||||
|
||||
context.getMqttService().subscribe(
|
||||
config.getProperty("mqtt.topic"),
|
||||
payload -> context.getBinaryEventService()
|
||||
.handlePayload(payload)
|
||||
);
|
||||
|
||||
context.setProcessManagerService(new ProcessManagerService(config));
|
||||
context.getProcessManagerService().startProcesses();
|
||||
|
||||
return context;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package vassistent.bootstrap;
|
||||
|
||||
public class ApplicationShutdownManager {
|
||||
|
||||
private final ApplicationContext context;
|
||||
|
||||
public ApplicationShutdownManager(ApplicationContext context) {
|
||||
this.context = context;
|
||||
}
|
||||
|
||||
public void shutdown() {
|
||||
|
||||
if (context.getMqttService() != null) {
|
||||
context.getMqttService().disconnect();
|
||||
}
|
||||
|
||||
if (context.getUnrealService() != null) {
|
||||
context.getUnrealService().disconnect();
|
||||
}
|
||||
|
||||
if (context.getProcessManagerService() != null) {
|
||||
context.getProcessManagerService().shutdown();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,97 @@
|
||||
package vassistent.controller;
|
||||
|
||||
import vassistent.bootstrap.ApplicationContext;
|
||||
import vassistent.bootstrap.ApplicationShutdownManager;
|
||||
import vassistent.model.AppState;
|
||||
import vassistent.ui.AppWindow;
|
||||
import vassistent.ui.DashboardView;
|
||||
import vassistent.util.Logger;
|
||||
|
||||
import java.awt.event.WindowAdapter;
|
||||
import java.awt.event.WindowEvent;
|
||||
|
||||
public class AppWindowController {
|
||||
|
||||
private final ApplicationContext context;
|
||||
private final ApplicationShutdownManager shutdownManager;
|
||||
|
||||
private AppWindow window;
|
||||
|
||||
public AppWindowController(ApplicationContext context) {
|
||||
this.context = context;
|
||||
this.shutdownManager = new ApplicationShutdownManager(context);
|
||||
subscribeAppState();
|
||||
}
|
||||
|
||||
public void createAndShowWindow() {
|
||||
|
||||
window = new AppWindow();
|
||||
|
||||
DashboardView dashboardView = window.getDashboardView();
|
||||
|
||||
DashboardController dashboardController =
|
||||
new DashboardController(
|
||||
context.getStatisticsService(),
|
||||
dashboardView,
|
||||
context.getAppState()
|
||||
);
|
||||
|
||||
window.getRefreshDashboardButton()
|
||||
.addActionListener(e -> dashboardController.loadChartData());
|
||||
|
||||
window.addWindowListener(new WindowAdapter() {
|
||||
@Override
|
||||
public void windowClosing(WindowEvent e) {
|
||||
shutdown();
|
||||
}
|
||||
});
|
||||
|
||||
window.setVisible(true);
|
||||
}
|
||||
|
||||
private void shutdown() {
|
||||
|
||||
shutdownManager.shutdown();
|
||||
|
||||
if (window != null)
|
||||
window.dispose();
|
||||
|
||||
System.exit(0);
|
||||
}
|
||||
|
||||
private void subscribeAppState() {
|
||||
|
||||
Logger.info("Controller",
|
||||
"Subscribe AppState Observer gestartet");
|
||||
|
||||
AppState state = context.getAppState();
|
||||
|
||||
state.addListener(appState -> {
|
||||
|
||||
Logger.debug("Controller",
|
||||
"AppState Update erhalten: " +
|
||||
appState.getProblemLevel());
|
||||
|
||||
if (window == null) {
|
||||
Logger.warn("Controller",
|
||||
"Window ist null, UI Update übersprungen");
|
||||
return;
|
||||
}
|
||||
|
||||
window.updateProblemLevel(
|
||||
appState.getProblemLevel().name()
|
||||
);
|
||||
|
||||
window.updateMqttStatus(
|
||||
appState.isMqttConnected()
|
||||
);
|
||||
|
||||
Logger.debug("Controller",
|
||||
"ProblemLevel UI aktualisiert → " +
|
||||
appState.getProblemLevel());
|
||||
});
|
||||
|
||||
Logger.info("Controller",
|
||||
"AppState Observer registriert");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package vassistent.controller;
|
||||
|
||||
import vassistent.model.AppState;
|
||||
import vassistent.model.RatioPoint;
|
||||
import vassistent.service.StatisticsService;
|
||||
import vassistent.ui.DashboardView;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
public class DashboardController {
|
||||
|
||||
private final StatisticsService statisticsService;
|
||||
private final DashboardView dashboardView;
|
||||
private long lastDataVersion = -1;
|
||||
|
||||
public DashboardController(
|
||||
StatisticsService statisticsService,
|
||||
DashboardView dashboardView,
|
||||
AppState appState
|
||||
) {
|
||||
this.statisticsService = statisticsService;
|
||||
this.dashboardView = dashboardView;
|
||||
|
||||
appState.addListener(state -> {
|
||||
|
||||
if (state.getDataVersion() != lastDataVersion) {
|
||||
lastDataVersion = state.getDataVersion();
|
||||
|
||||
List<RatioPoint> points = statisticsService.getLastNAverages(20);
|
||||
dashboardView.updateChart(points);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public void loadChartData() {
|
||||
|
||||
List<RatioPoint> points = statisticsService.getLastNAverages(20);
|
||||
|
||||
dashboardView.updateChart(points);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
package vassistent.controller;
|
||||
|
||||
import vassistent.ui.PixelStreamingView;
|
||||
|
||||
public class StreamingController {
|
||||
|
||||
private final PixelStreamingView view;
|
||||
|
||||
public StreamingController(PixelStreamingView view) {
|
||||
this.view = view;
|
||||
}
|
||||
|
||||
public void reloadStream(String url) {
|
||||
view.loadStream(url);
|
||||
}
|
||||
|
||||
public void toggleFullscreen() {
|
||||
view.requestFocus();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
package vassistent.model;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public class AppState {
|
||||
|
||||
private long dataVersion = 0;
|
||||
private boolean mqttConnected;
|
||||
private ProblemLevel problemLevel = ProblemLevel.NONE;
|
||||
|
||||
private final List<Consumer<AppState>> listeners =
|
||||
new ArrayList<>();
|
||||
|
||||
public void addListener(Consumer<AppState> listener) {
|
||||
listeners.add(listener);
|
||||
}
|
||||
|
||||
private void notifyListeners() {
|
||||
listeners.forEach(l -> l.accept(this));
|
||||
}
|
||||
|
||||
public ProblemLevel getProblemLevel() {
|
||||
return problemLevel;
|
||||
}
|
||||
|
||||
public void setProblemLevel(ProblemLevel problemLevel) {
|
||||
this.problemLevel = problemLevel;
|
||||
notifyListeners();
|
||||
}
|
||||
|
||||
public boolean isMqttConnected() {
|
||||
return mqttConnected;
|
||||
}
|
||||
|
||||
public void setMqttConnected(boolean mqttConnected) {
|
||||
if (this.mqttConnected != mqttConnected) {
|
||||
this.mqttConnected = mqttConnected;
|
||||
notifyListeners();
|
||||
}
|
||||
}
|
||||
|
||||
public long getDataVersion() {
|
||||
return dataVersion;
|
||||
}
|
||||
|
||||
public void incrementDataVersion() {
|
||||
dataVersion++;
|
||||
notifyListeners();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
package vassistent.model;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
public class DatabaseEntry {
|
||||
|
||||
private final LocalDateTime timestamp;
|
||||
private final double value;
|
||||
|
||||
|
||||
public DatabaseEntry(LocalDateTime timestamp, double value) {
|
||||
this.timestamp = timestamp;
|
||||
this.value = value;
|
||||
}
|
||||
|
||||
public LocalDateTime getTimestamp() {
|
||||
return timestamp;
|
||||
}
|
||||
|
||||
public double getValue() {
|
||||
return value;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
package vassistent.model;
|
||||
|
||||
public enum ProblemLevel {
|
||||
NONE,
|
||||
WARNING,
|
||||
HIGH,
|
||||
DISASTER
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
package vassistent.model;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
public class RatioPoint {
|
||||
|
||||
private final LocalDateTime timestamp;
|
||||
private final double ratio;
|
||||
|
||||
|
||||
public RatioPoint(LocalDateTime timestamp, double ratio) {
|
||||
this.timestamp = timestamp;
|
||||
this.ratio = ratio;
|
||||
}
|
||||
|
||||
public LocalDateTime getTimestamp() {
|
||||
return timestamp;
|
||||
}
|
||||
|
||||
public double getRatio() {
|
||||
return ratio;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
package vassistent.service;
|
||||
|
||||
import vassistent.util.Logger;
|
||||
|
||||
public class BinaryEventService {
|
||||
private final DataPersistenceService persistenceService;
|
||||
private final EvaluationService evaluationService;
|
||||
|
||||
public BinaryEventService(
|
||||
DataPersistenceService persistenceService,
|
||||
EvaluationService evaluationService
|
||||
) {
|
||||
this.persistenceService = persistenceService;
|
||||
this.evaluationService = evaluationService;
|
||||
}
|
||||
|
||||
public void handlePayload(String payload) {
|
||||
try {
|
||||
int value = Integer.parseInt(payload);
|
||||
|
||||
if (value != 0 && value != 1) {
|
||||
Logger.warn("EVENT", "Ungültiger Wert: " + payload);
|
||||
return;
|
||||
}
|
||||
persistenceService.store(value);
|
||||
evaluationService.evaluate();
|
||||
|
||||
} catch (NumberFormatException e) {
|
||||
Logger.warn("EVENT", "Payload nicht numerisch " + payload);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,123 @@
|
||||
package vassistent.service;
|
||||
|
||||
import vassistent.model.DatabaseEntry;
|
||||
import vassistent.model.RatioPoint;
|
||||
import vassistent.util.Logger;
|
||||
|
||||
import java.io.File;
|
||||
import java.sql.*;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
|
||||
public class DataPersistenceService {
|
||||
|
||||
private static final String DB_FOLDER = "data";
|
||||
private final String DB_URL;
|
||||
|
||||
public DataPersistenceService() {
|
||||
this("jdbc:sqlite:data/health.db");
|
||||
createDatabaseFolder();
|
||||
init();
|
||||
}
|
||||
|
||||
// Für Tests
|
||||
public DataPersistenceService(String DB_URL) {
|
||||
this.DB_URL = DB_URL;
|
||||
createDatabaseFolder();
|
||||
init();
|
||||
}
|
||||
|
||||
public void createDatabaseFolder() {
|
||||
File folder = new File(DB_FOLDER);
|
||||
if (!folder.exists()) {
|
||||
boolean created = folder.mkdirs();
|
||||
if (created) {
|
||||
Logger.info("DB", "Datenbank-Ordner erstellt");
|
||||
} else {
|
||||
Logger.warn("DB", "Konnte DB-Ordner nicht erstellen");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void init() {
|
||||
try {
|
||||
Class.forName("org.sqlite.JDBC");
|
||||
} catch (ClassNotFoundException e) {
|
||||
Logger.error("DB", "SQLite JDBC Driver nicht gefunden", e);
|
||||
}
|
||||
|
||||
try (
|
||||
Connection c = DriverManager.getConnection(DB_URL); Statement s = c.createStatement()) {
|
||||
|
||||
s.execute("""
|
||||
CREATE TABLE IF NOT EXISTS binary_event (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
value INTEGER NOT NULL,
|
||||
timestamp DATETIME NOT NULL
|
||||
)
|
||||
""");
|
||||
|
||||
Logger.info("DB", "Datenbank initialisiert");
|
||||
|
||||
} catch (SQLException e) {
|
||||
Logger.error("DB", "Datenbank initialisierung fehlgeschlagen", e);
|
||||
}
|
||||
}
|
||||
|
||||
public void store(int value) {
|
||||
try (
|
||||
Connection c = DriverManager.getConnection(DB_URL); PreparedStatement ps = c.prepareStatement(
|
||||
"INSERT INTO binary_event (value, timestamp) VALUES (?, ?)"
|
||||
)) {
|
||||
ps.setInt(1, value);
|
||||
ps.setString(2, LocalDateTime.now().toString());
|
||||
ps.executeUpdate();
|
||||
Logger.debug("DB", "Database Insert erfolgreich: " + ps.toString());
|
||||
} catch (SQLException e) {
|
||||
Logger.error("DB", "Database Insert fehlgeschlagen", e);
|
||||
}
|
||||
}
|
||||
|
||||
public List<DatabaseEntry> getLastEntries(int count) {
|
||||
List<DatabaseEntry> result = new ArrayList<>();
|
||||
|
||||
String query = """
|
||||
SELECT value, timestamp
|
||||
FROM binary_event
|
||||
ORDER BY timestamp DESC
|
||||
LIMIT ?
|
||||
""";
|
||||
try (Connection c = DriverManager.getConnection(DB_URL); PreparedStatement ps = c.prepareStatement(query)) {
|
||||
|
||||
ps.setInt(1, count);
|
||||
|
||||
ResultSet rs = ps.executeQuery();
|
||||
|
||||
while (rs.next()) {
|
||||
|
||||
String ts = rs.getString("timestamp");
|
||||
|
||||
LocalDateTime timestamp =
|
||||
LocalDateTime.parse(ts);
|
||||
|
||||
double value = rs.getDouble("value");
|
||||
|
||||
result.add(new DatabaseEntry(timestamp, value));
|
||||
}
|
||||
|
||||
} catch (SQLException e) {
|
||||
Logger.error("DB",
|
||||
"Fehler beim Laden der letzten Einträge", e);
|
||||
}
|
||||
|
||||
Collections.reverse(result);
|
||||
return result;
|
||||
}
|
||||
|
||||
public String getDbUrl() {
|
||||
return DB_URL;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
package vassistent.service;
|
||||
|
||||
|
||||
import vassistent.model.AppState;
|
||||
import vassistent.model.ProblemLevel;
|
||||
|
||||
public class EvaluationService {
|
||||
private final StatisticsService statisticsService;
|
||||
private final AppState appState;
|
||||
private final UnrealWebSocketService unrealService;
|
||||
|
||||
public EvaluationService(
|
||||
StatisticsService statisticsService,
|
||||
AppState appState,
|
||||
UnrealWebSocketService unrealService
|
||||
) {
|
||||
|
||||
this.statisticsService = statisticsService;
|
||||
this.appState = appState;
|
||||
this.unrealService = unrealService;
|
||||
}
|
||||
|
||||
public void evaluate() {
|
||||
double ratio = statisticsService.getRatio(10);
|
||||
|
||||
ProblemLevel level = calculateLevel(ratio);
|
||||
|
||||
appState.setProblemLevel(level);
|
||||
unrealService.speak(level.name());
|
||||
|
||||
appState.incrementDataVersion();
|
||||
}
|
||||
|
||||
private ProblemLevel calculateLevel(Double ratio) {
|
||||
if (ratio >= 0.9) {
|
||||
return ProblemLevel.DISASTER;
|
||||
} else if (ratio >= 0.8) {
|
||||
return ProblemLevel.HIGH;
|
||||
} else if (ratio >= 0.5) {
|
||||
return ProblemLevel.WARNING;
|
||||
} else {
|
||||
return ProblemLevel.NONE;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,127 @@
|
||||
package vassistent.service;
|
||||
|
||||
import org.eclipse.paho.client.mqttv3.*;
|
||||
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
|
||||
import vassistent.model.AppState;
|
||||
import vassistent.util.Logger;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public class MqttClientService implements MqttCallback {
|
||||
|
||||
private AppState appState;
|
||||
private static final String BROKER_URL = "tcp://localhost:1883";
|
||||
private static final String CLIENT_ID = "JavaClientPublisherSubscriber";
|
||||
|
||||
private final Map<String, Consumer<String>> topicListeners =
|
||||
new ConcurrentHashMap<>();
|
||||
|
||||
private MqttClient client;
|
||||
|
||||
public MqttClientService(AppState appState) {
|
||||
this.appState = appState;
|
||||
|
||||
try {
|
||||
client = new MqttClient(
|
||||
BROKER_URL,
|
||||
CLIENT_ID,
|
||||
new MemoryPersistence()
|
||||
);
|
||||
client.setCallback(this);
|
||||
|
||||
MqttConnectOptions options = new MqttConnectOptions();
|
||||
options.setCleanSession(true);
|
||||
options.setAutomaticReconnect(true);
|
||||
|
||||
Logger.info("MQTT", "Verbinde mit Broker " + BROKER_URL);
|
||||
client.connect(options);
|
||||
appState.setMqttConnected(true);
|
||||
Logger.info("MQTT", "Verbindung hergestellt");
|
||||
|
||||
} catch (MqttException e) {
|
||||
Logger.error("MQTT", "Fehler beim Verbinden", e);
|
||||
}
|
||||
}
|
||||
|
||||
public void disconnect() {
|
||||
try {
|
||||
if (client != null && client.isConnected()) {
|
||||
client.disconnect();
|
||||
Logger.info("MQTT", "Verbindung getrennt");
|
||||
}
|
||||
} catch (MqttException e) {
|
||||
Logger.error("MQTT", "Fehler beim Trennen", e);
|
||||
}
|
||||
}
|
||||
|
||||
public void subscribe(String topic, Consumer<String> listener) {
|
||||
topicListeners.put(topic, listener);
|
||||
try {
|
||||
client.subscribe(topic);
|
||||
Logger.info("MQTT", "Topic abonniert: " + topic);
|
||||
} catch (MqttException e) {
|
||||
Logger.error("MQTT", "Subscribe fehlgeschlagen: " + topic, e);
|
||||
}
|
||||
}
|
||||
|
||||
public void publish(String topic, String message, int qos) {
|
||||
try {
|
||||
MqttMessage mqttMessage =
|
||||
new MqttMessage(message.getBytes(StandardCharsets.UTF_8));
|
||||
mqttMessage.setQos(qos);
|
||||
|
||||
client.publish(topic, mqttMessage);
|
||||
|
||||
Logger.debug(
|
||||
"MQTT",
|
||||
"Publish -> Topic=" + topic + ", QoS=" + qos + ", Payload=" + message
|
||||
);
|
||||
} catch (MqttException e) {
|
||||
Logger.error("MQTT", "Publish fehlgeschlagen", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void connectionLost(Throwable cause) {
|
||||
Logger.warn(
|
||||
"MQTT",
|
||||
"Verbindung verloren: " + cause.getMessage()
|
||||
);
|
||||
|
||||
appState.setMqttConnected(false);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void messageArrived(String topic, MqttMessage message) throws Exception {
|
||||
String payload =
|
||||
new String(message.getPayload(), StandardCharsets.UTF_8);
|
||||
|
||||
Logger.debug(
|
||||
"MQTT",
|
||||
"Nachricht empfangen -> Topic=" + topic + ", Payload=" + payload
|
||||
);
|
||||
|
||||
Consumer<String> listener = topicListeners.get(topic);
|
||||
if (listener != null) {
|
||||
listener.accept(payload);
|
||||
} else {
|
||||
Logger.warn(
|
||||
"MQTT",
|
||||
"Keine Listener für Topic: " + topic
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deliveryComplete(IMqttDeliveryToken token) {
|
||||
Logger.debug(
|
||||
"MQTT",
|
||||
"Delivery complete, MessageId=" + token.getMessageId()
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
package vassistent.service;
|
||||
|
||||
public class PixelStreamingService {
|
||||
|
||||
}
|
||||
@@ -0,0 +1,79 @@
|
||||
package vassistent.service;
|
||||
|
||||
import vassistent.util.Logger;
|
||||
|
||||
import java.io.File;
|
||||
import java.io.IOException;
|
||||
import java.util.Properties;
|
||||
|
||||
public class ProcessManagerService {
|
||||
|
||||
private final Properties config;
|
||||
|
||||
private Process pythonProcess;
|
||||
private Process unrealProcess;
|
||||
|
||||
public ProcessManagerService(Properties config) {
|
||||
this.config = config;
|
||||
}
|
||||
|
||||
public void startProcesses() {
|
||||
|
||||
startPythonIfEnabled();
|
||||
startUnrealIfEnabled();
|
||||
}
|
||||
|
||||
private void startPythonIfEnabled() {
|
||||
|
||||
boolean enabled = Boolean.parseBoolean(
|
||||
config.getProperty("mqtt_sim.enabled", "false")
|
||||
);
|
||||
|
||||
if (!enabled) return;
|
||||
|
||||
try {
|
||||
String script = config.getProperty("mqtt_sim.script");
|
||||
|
||||
ProcessBuilder pb = new ProcessBuilder("python", script);
|
||||
|
||||
pb.directory(new File("."));
|
||||
pythonProcess = pb.start();
|
||||
|
||||
Logger.info("PROCESS", "Mqtt Simulator gestartet");
|
||||
} catch (IOException e) {
|
||||
Logger.error("PROCESS", "Mqtt Simulator Start fehlgeschlagen", e);
|
||||
}
|
||||
}
|
||||
|
||||
private void startUnrealIfEnabled() {
|
||||
|
||||
boolean enabled = Boolean.parseBoolean(config.getProperty("unreal.enabled", "false"));
|
||||
|
||||
if (!enabled) return;
|
||||
|
||||
try {
|
||||
|
||||
String exe = config.getProperty("unreal.executable");
|
||||
|
||||
ProcessBuilder pb = new ProcessBuilder(exe);
|
||||
|
||||
unrealProcess = pb.start();
|
||||
|
||||
Logger.info("PROCESS", "Unreal Engine Avatar");
|
||||
|
||||
} catch (IOException e) {
|
||||
Logger.error("PROCESS", "Unreal Engine Start fehlgeschlagen", e);
|
||||
}
|
||||
}
|
||||
|
||||
public void shutdown() {
|
||||
|
||||
if (pythonProcess != null)
|
||||
pythonProcess.destroy();
|
||||
|
||||
if (unrealProcess != null)
|
||||
unrealProcess.destroy();
|
||||
|
||||
Logger.info("PROCESS", "Externe Prozesse beendet");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
package vassistent.service;
|
||||
|
||||
import vassistent.model.AppState;
|
||||
|
||||
public class StateService {
|
||||
private static final AppState STATE = new AppState();
|
||||
|
||||
private StateService() {}
|
||||
|
||||
public static AppState getState() {
|
||||
return STATE;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,69 @@
|
||||
package vassistent.service;
|
||||
|
||||
import vassistent.model.DatabaseEntry;
|
||||
import vassistent.model.RatioPoint;
|
||||
import vassistent.util.Logger;
|
||||
|
||||
import java.sql.*;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
public class StatisticsService {
|
||||
private final DataPersistenceService persistenceService;
|
||||
|
||||
public StatisticsService(DataPersistenceService persistenceService) {
|
||||
this.persistenceService = persistenceService;
|
||||
}
|
||||
|
||||
public double getRatio(int lastN) {
|
||||
String query = """
|
||||
SELECT AVG(value)
|
||||
FROM (
|
||||
SELECT value
|
||||
FROM binary_event
|
||||
ORDER BY id DESC
|
||||
LIMIT ?
|
||||
)
|
||||
""";
|
||||
|
||||
try (Connection c = DriverManager.getConnection(persistenceService.getDbUrl()); PreparedStatement ps = c.prepareStatement(query)) {
|
||||
ps.setInt(1, lastN);
|
||||
|
||||
ResultSet rs = ps.executeQuery();
|
||||
if (rs.next()) {
|
||||
return rs.getDouble(1);
|
||||
}
|
||||
} catch (SQLException e) {
|
||||
Logger.error("Evaluation", "Couldn't get ratio.", e);
|
||||
}
|
||||
|
||||
return 0.0;
|
||||
}
|
||||
|
||||
public List<RatioPoint> getLastNAverages(int lastN) {
|
||||
|
||||
List<DatabaseEntry> entries = persistenceService.getLastEntries(30);
|
||||
|
||||
List<RatioPoint> result = new ArrayList<>();
|
||||
|
||||
if (entries.isEmpty())
|
||||
return result;
|
||||
|
||||
for (int i = 10; i < entries.size(); i++) {
|
||||
double sum = 0;
|
||||
|
||||
for (int j = i - 10; j < i; j++) {
|
||||
sum += entries.get(j).getValue();
|
||||
}
|
||||
|
||||
double avg = sum / 10.0;
|
||||
|
||||
LocalDateTime timestamp = entries.get(i).getTimestamp();
|
||||
|
||||
result.add(new RatioPoint(timestamp, avg));
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,95 @@
|
||||
package vassistent.service;
|
||||
|
||||
import vassistent.util.Logger;
|
||||
|
||||
import java.net.URI;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import jakarta.websocket.*;
|
||||
|
||||
public class UnrealWebSocketService {
|
||||
|
||||
private Session session;
|
||||
private final URI serverUri;
|
||||
private final AtomicBoolean connected = new AtomicBoolean(false);
|
||||
|
||||
|
||||
public UnrealWebSocketService(String serverUrl) {
|
||||
this.serverUri = URI.create(serverUrl);
|
||||
connect();
|
||||
}
|
||||
|
||||
private void connect() {
|
||||
try {
|
||||
WebSocketContainer container =
|
||||
ContainerProvider.getWebSocketContainer();
|
||||
|
||||
Logger.info("UNREAL", "Verbinde zu " + serverUri);
|
||||
container.connectToServer(this, serverUri);
|
||||
|
||||
} catch (Exception e) {
|
||||
Logger.error("UNREAL", "WebSocket-Verbindung fehlgeschlagen", e);
|
||||
}
|
||||
}
|
||||
|
||||
public boolean isConnected() {
|
||||
return connected.get();
|
||||
}
|
||||
|
||||
public void disconnect() {
|
||||
try {
|
||||
if (session != null && session.isOpen()) {
|
||||
session.close();
|
||||
Logger.info("UNREAL", "WebSocket-Verbindung geschlossen");
|
||||
}
|
||||
} catch (Exception e) {
|
||||
Logger.error("UNREAL", "Fehler beim Schließen der Verbindung", e);
|
||||
}
|
||||
}
|
||||
|
||||
@OnOpen
|
||||
public void onOpen(Session session) {
|
||||
this.session = session;
|
||||
connected.set(true);
|
||||
Logger.info("UNREAL", "WebSocket verbunden");
|
||||
}
|
||||
|
||||
@OnClose
|
||||
public void onClose(Session session, CloseReason reason) {
|
||||
connected.set(false);
|
||||
Logger.warn("UNREAL", "WebSocket geschlossen: " + reason.getReasonPhrase());
|
||||
}
|
||||
|
||||
@OnError
|
||||
public void onError(Session session, Throwable throwable) {
|
||||
connected.set(false);
|
||||
Logger.error("UNREAL", "WebSocket-Fehler", throwable);
|
||||
}
|
||||
|
||||
@OnMessage
|
||||
public void onMessage(String message) {
|
||||
Logger.debug("UNREAL", "Nachricht empfangen: " + message);
|
||||
}
|
||||
|
||||
public void speak(String text) {
|
||||
if (!connected.get()) {
|
||||
Logger.warn("UNREAL", "Nicht verbunden! Text wird nicht gesendet!");
|
||||
return;
|
||||
}
|
||||
|
||||
String json = buildSpeakMessage(text);
|
||||
|
||||
session.getAsyncRemote().sendText(json);
|
||||
Logger.info("UNREAL", "Sende Speak-Command: " + text);
|
||||
}
|
||||
|
||||
public String buildSpeakMessage(String text) {
|
||||
return "{"
|
||||
+ "\"type\":\"speak\","
|
||||
+ "\"text\":\"" + escape(text) + "\""
|
||||
+ "}";
|
||||
}
|
||||
|
||||
public String escape(String text) {
|
||||
return text.replace("\"", "\\\"");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,125 @@
|
||||
package vassistent.ui;
|
||||
|
||||
import vassistent.bootstrap.ApplicationContext;
|
||||
import vassistent.controller.DashboardController;
|
||||
|
||||
import javax.swing.*;
|
||||
import java.awt.*;
|
||||
|
||||
public class AppWindow extends JFrame {
|
||||
|
||||
private PixelStreamingView streamingView;
|
||||
private DashboardView dashboardView;
|
||||
|
||||
private JLabel mqttStatusLabel;
|
||||
private JLabel websocketStatusLabel;
|
||||
private JLabel problemLevelLabel;
|
||||
|
||||
private JButton refreshDashboardButton;
|
||||
private JButton reloadStreamButton;
|
||||
private JButton fullscreenButton;
|
||||
|
||||
public AppWindow() {
|
||||
|
||||
setTitle("Virtueller Gesundheitsassistent");
|
||||
setDefaultCloseOperation(JFrame.EXIT_ON_CLOSE);
|
||||
setSize(1400, 850);
|
||||
setLocationRelativeTo(null);
|
||||
|
||||
setLayout(new BorderLayout());
|
||||
|
||||
add(createToolbar(), BorderLayout.NORTH);
|
||||
add(createTabPane(), BorderLayout.CENTER);
|
||||
add(createStatusBar(), BorderLayout.SOUTH);
|
||||
}
|
||||
|
||||
// ---------------- Toolbar ----------------
|
||||
|
||||
private JPanel createToolbar() {
|
||||
|
||||
JPanel toolbar = new JPanel(new FlowLayout(FlowLayout.LEFT));
|
||||
|
||||
refreshDashboardButton = new JButton("Dashboard Refresh");
|
||||
reloadStreamButton = new JButton("Stream Reload");
|
||||
fullscreenButton = new JButton("Fullscreen Stream");
|
||||
|
||||
toolbar.add(refreshDashboardButton);
|
||||
toolbar.add(reloadStreamButton);
|
||||
toolbar.add(fullscreenButton);
|
||||
|
||||
return toolbar;
|
||||
}
|
||||
|
||||
// ---------------- Tabs ----------------
|
||||
|
||||
private JTabbedPane createTabPane() {
|
||||
|
||||
JTabbedPane tabs = new JTabbedPane();
|
||||
|
||||
streamingView = new PixelStreamingView(
|
||||
"http://google.com",
|
||||
false,
|
||||
false
|
||||
);
|
||||
|
||||
dashboardView = new DashboardView();
|
||||
|
||||
tabs.addTab("Avatar Streaming", streamingView);
|
||||
tabs.addTab("Dashboard", dashboardView);
|
||||
|
||||
return tabs;
|
||||
}
|
||||
|
||||
// ---------------- Status Bar ----------------
|
||||
|
||||
private JPanel createStatusBar() {
|
||||
|
||||
JPanel statusBar = new JPanel(new FlowLayout(FlowLayout.LEFT));
|
||||
|
||||
mqttStatusLabel = new JLabel("MQTT: Disconnected");
|
||||
websocketStatusLabel = new JLabel();
|
||||
problemLevelLabel = new JLabel("Problem: NONE");
|
||||
|
||||
statusBar.add(mqttStatusLabel);
|
||||
statusBar.add(new JLabel(" | "));
|
||||
statusBar.add(websocketStatusLabel);
|
||||
statusBar.add(new JLabel(" | "));
|
||||
statusBar.add(problemLevelLabel);
|
||||
|
||||
return statusBar;
|
||||
}
|
||||
|
||||
public void updateMqttStatus(boolean connected) {
|
||||
mqttStatusLabel.setText("MQTT: " +
|
||||
(connected ? "Connected" : "Disconnected"));
|
||||
}
|
||||
|
||||
public void updateWebsocketStatus(boolean connected) {
|
||||
websocketStatusLabel.setText("WebSocket: " +
|
||||
(connected ? "Connected" : "Disconnected"));
|
||||
}
|
||||
|
||||
public void updateProblemLevel(String level) {
|
||||
problemLevelLabel.setText("Problem: " + level);
|
||||
}
|
||||
|
||||
public DashboardView getDashboardView() {
|
||||
return dashboardView;
|
||||
}
|
||||
|
||||
public PixelStreamingView getStreamingView() {
|
||||
return streamingView;
|
||||
}
|
||||
|
||||
public JButton getRefreshDashboardButton() {
|
||||
return refreshDashboardButton;
|
||||
}
|
||||
|
||||
public JButton getReloadStreamButton() {
|
||||
return reloadStreamButton;
|
||||
}
|
||||
|
||||
public JButton getFullscreenButton() {
|
||||
return fullscreenButton;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,63 @@
|
||||
package vassistent.ui;
|
||||
|
||||
import org.jfree.chart.ChartFactory;
|
||||
import org.jfree.chart.ChartPanel;
|
||||
import org.jfree.chart.JFreeChart;
|
||||
import org.jfree.data.category.DefaultCategoryDataset;
|
||||
import org.jfree.data.time.Millisecond;
|
||||
import org.jfree.data.time.TimeSeries;
|
||||
import org.jfree.data.time.TimeSeriesCollection;
|
||||
import vassistent.model.RatioPoint;
|
||||
|
||||
import javax.swing.*;
|
||||
import java.awt.*;
|
||||
import java.sql.Timestamp;
|
||||
import java.time.ZoneId;
|
||||
import java.util.Arrays;
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
|
||||
public class DashboardView extends JPanel {
|
||||
|
||||
private TimeSeries series;
|
||||
|
||||
public DashboardView() {
|
||||
setLayout(new BorderLayout());
|
||||
|
||||
series = new TimeSeries("Ratio");
|
||||
|
||||
TimeSeriesCollection dataset =
|
||||
new TimeSeriesCollection(series);
|
||||
|
||||
JFreeChart chart = ChartFactory.createTimeSeriesChart(
|
||||
"Ratio Verlauf",
|
||||
"Zeit",
|
||||
"Durchschnitt",
|
||||
dataset
|
||||
);
|
||||
|
||||
ChartPanel chartPanel = new ChartPanel(chart);
|
||||
|
||||
add(chartPanel, BorderLayout.CENTER);
|
||||
}
|
||||
|
||||
public void updateChart(List<RatioPoint> points) {
|
||||
series.clear();
|
||||
|
||||
for (RatioPoint point : points) {
|
||||
|
||||
series.addOrUpdate(
|
||||
new Millisecond(
|
||||
Date.from(
|
||||
point.getTimestamp()
|
||||
.atZone(ZoneId.systemDefault())
|
||||
.toInstant()
|
||||
)
|
||||
),
|
||||
point.getRatio()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
package vassistent.ui;
|
||||
|
||||
import org.cef.CefApp;
|
||||
import org.cef.CefClient;
|
||||
import org.cef.CefSettings;
|
||||
import org.cef.browser.CefBrowser;
|
||||
import org.cef.handler.CefAppHandlerAdapter;
|
||||
|
||||
import javax.swing.*;
|
||||
import java.awt.*;
|
||||
import java.awt.event.ActionEvent;
|
||||
import java.awt.event.ActionListener;
|
||||
|
||||
public class PixelStreamingView extends JPanel {
|
||||
|
||||
private static final long serialVersionUID = -5570653778104813836L;
|
||||
private final CefApp cefApp;
|
||||
private final CefClient client;
|
||||
private final CefBrowser browser;
|
||||
private final Component browserUI_;
|
||||
|
||||
public PixelStreamingView(String startURL, boolean useOSR, boolean isTransparent) {
|
||||
super(new BorderLayout());
|
||||
|
||||
CefApp.addAppHandler(new CefAppHandlerAdapter(null) {
|
||||
@Override
|
||||
public void stateHasChanged(CefApp.CefAppState state) {
|
||||
if (state == CefApp.CefAppState.TERMINATED){
|
||||
System.out.println("TERMINATED");
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
CefSettings settings = new CefSettings();
|
||||
settings.windowless_rendering_enabled = useOSR;
|
||||
|
||||
cefApp = CefApp.getInstance(settings);
|
||||
client = cefApp.createClient();
|
||||
|
||||
browser = client.createBrowser(startURL, useOSR, isTransparent);
|
||||
browserUI_ = browser.getUIComponent();
|
||||
|
||||
add(browserUI_, BorderLayout.CENTER);
|
||||
}
|
||||
|
||||
public void dispose() {
|
||||
if (browser != null) browser.close(true);
|
||||
if (client != null) client.dispose();
|
||||
if (cefApp != null) cefApp.dispose();
|
||||
}
|
||||
|
||||
public void loadStream(String url) {
|
||||
browser.loadURL(url);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
package vassistent.util;
|
||||
|
||||
import java.io.InputStream;
|
||||
import java.util.Properties;
|
||||
|
||||
public class ConfigLoader {
|
||||
|
||||
public ConfigLoader() {}
|
||||
|
||||
public static Properties loadProperties(String path) {
|
||||
|
||||
Properties props = new Properties();
|
||||
|
||||
try (InputStream is =
|
||||
Logger.class.getClassLoader()
|
||||
.getResourceAsStream(path)) {
|
||||
|
||||
if (is != null)
|
||||
props.load(is);
|
||||
|
||||
} catch (Exception e) {
|
||||
Logger.error("CONFIG",
|
||||
"Properties Laden fehlgeschlagen", e);
|
||||
}
|
||||
|
||||
return props;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,199 @@
|
||||
package vassistent.util;
|
||||
|
||||
import java.io.*;
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.format.DateTimeFormatter;
|
||||
import java.util.Properties;
|
||||
import java.util.concurrent.BlockingQueue;
|
||||
import java.util.concurrent.LinkedBlockingQueue;
|
||||
|
||||
public class Logger {
|
||||
|
||||
public enum Level {
|
||||
DEBUG, INFO, WARN, ERROR
|
||||
}
|
||||
|
||||
private static Level currentLevel = Level.DEBUG;
|
||||
|
||||
private static boolean logToFile = true;
|
||||
private static String logFilePath = "logs/application.log";
|
||||
private static int maxFileSizeMB = 10;
|
||||
|
||||
private static final DateTimeFormatter FORMATTER =
|
||||
DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
|
||||
|
||||
private static final BlockingQueue<String> logQueue =
|
||||
new LinkedBlockingQueue<>();
|
||||
|
||||
static {
|
||||
loadConfig();
|
||||
startLogWriterThread();
|
||||
}
|
||||
|
||||
private Logger() {}
|
||||
|
||||
// ================= CONFIG =================
|
||||
|
||||
private static void loadConfig() {
|
||||
try (InputStream is = Logger.class.getClassLoader()
|
||||
.getResourceAsStream("config/logger.properties")) {
|
||||
|
||||
if (is == null) return;
|
||||
|
||||
Properties props = new Properties();
|
||||
props.load(is);
|
||||
|
||||
currentLevel =
|
||||
Level.valueOf(
|
||||
props.getProperty("logger.level", "DEBUG"));
|
||||
|
||||
logToFile =
|
||||
Boolean.parseBoolean(
|
||||
props.getProperty("logger.file.enabled", "true"));
|
||||
|
||||
logFilePath =
|
||||
props.getProperty("logger.file", "logs/application.log");
|
||||
|
||||
maxFileSizeMB =
|
||||
Integer.parseInt(
|
||||
props.getProperty("logger.max.size.mb", "10"));
|
||||
|
||||
} catch (Exception e) {
|
||||
System.err.println("Logger Config Fehler: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
// ================= API =================
|
||||
|
||||
public static void debug(String source, String message) {
|
||||
log(Level.DEBUG, source, message, null);
|
||||
}
|
||||
|
||||
public static void info(String source, String message) {
|
||||
log(Level.INFO, source, message, null);
|
||||
}
|
||||
|
||||
public static void warn(String source, String message) {
|
||||
log(Level.WARN, source, message, null);
|
||||
}
|
||||
|
||||
public static void error(String source, String message) {
|
||||
log(Level.ERROR, source, message, null);
|
||||
}
|
||||
|
||||
public static void error(String source, String message, Throwable t) {
|
||||
log(Level.ERROR, source, message, t);
|
||||
}
|
||||
|
||||
// ================= CORE =================
|
||||
|
||||
public static void log(Level level, String source, String message, Throwable throwable) {
|
||||
|
||||
if (level.ordinal() < currentLevel.ordinal()) return;
|
||||
|
||||
String timestamp =
|
||||
LocalDateTime.now().format(FORMATTER);
|
||||
|
||||
StringBuilder sb = new StringBuilder();
|
||||
|
||||
sb.append(String.format("[%s] [%s] [%s] %s",
|
||||
timestamp, level, source, message));
|
||||
|
||||
if (throwable != null) {
|
||||
StringWriter sw = new StringWriter();
|
||||
throwable.printStackTrace(new PrintWriter(sw));
|
||||
sb.append("\n").append(sw);
|
||||
}
|
||||
|
||||
logQueue.offer(sb.toString());
|
||||
}
|
||||
|
||||
// ================= BACKGROUND WRITER =================
|
||||
|
||||
private static void startLogWriterThread() {
|
||||
|
||||
Thread writerThread = new Thread(() -> {
|
||||
|
||||
while (true) {
|
||||
|
||||
try {
|
||||
|
||||
String logEntry = logQueue.take();
|
||||
|
||||
if (logToFile) {
|
||||
|
||||
rotateIfNeeded();
|
||||
|
||||
File logFile = new File(logFilePath);
|
||||
logFile.getParentFile().mkdirs();
|
||||
|
||||
try (FileWriter fw = new FileWriter(logFile, true)) {
|
||||
fw.write(logEntry + "\n");
|
||||
}
|
||||
}
|
||||
|
||||
System.out.println(logEntry);
|
||||
|
||||
} catch (Exception e) {
|
||||
System.err.println("Logger Writer Fehler: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
});
|
||||
|
||||
writerThread.setDaemon(true);
|
||||
writerThread.start();
|
||||
}
|
||||
|
||||
// ================= ROTATION =================
|
||||
|
||||
private static void rotateIfNeeded() {
|
||||
|
||||
try {
|
||||
|
||||
File file = new File(logFilePath);
|
||||
|
||||
if (!file.exists()) return;
|
||||
|
||||
long maxBytes =
|
||||
maxFileSizeMB * 1024L * 1024L;
|
||||
|
||||
if (file.length() > maxBytes) {
|
||||
|
||||
String archiveName =
|
||||
logFilePath.replace(".log",
|
||||
"_" + System.currentTimeMillis() + ".log");
|
||||
|
||||
File archive = new File(archiveName);
|
||||
|
||||
boolean renamed = file.renameTo(archive);
|
||||
|
||||
if (renamed) {
|
||||
System.out.println("Logger Rotation: " + archiveName);
|
||||
}
|
||||
}
|
||||
|
||||
} catch (Exception e) {
|
||||
System.err.println("Logger Rotation Fehler: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
public static void shutdown() {
|
||||
|
||||
System.out.println("Logger Shutdown gestartet");
|
||||
|
||||
try {
|
||||
Thread.sleep(500); // Warte auf Queue Flush
|
||||
|
||||
while (!logQueue.isEmpty()) {
|
||||
String log = logQueue.poll();
|
||||
if (log != null) {
|
||||
System.out.println(log);
|
||||
}
|
||||
}
|
||||
|
||||
} catch (Exception e) {
|
||||
System.err.println("Logger Shutdown Fehler: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
# ===== MODE =====
|
||||
app.mode=test
|
||||
|
||||
# ===== MQTT CLIENT =====
|
||||
mqtt.topic="PREDICTION"
|
||||
|
||||
# ===== MQTT SIMULATOR =====
|
||||
mqtt_sim.enabled=true
|
||||
mqtt_sim.script=scripts/mqtt_simulator.py
|
||||
|
||||
# ===== UNREAL ENGINE =====
|
||||
unreal.enabled=true
|
||||
unreal.executable=external/unreal/avatar.exe
|
||||
unreal.signalling_server.script=external/unreal/start.bat
|
||||
@@ -0,0 +1,4 @@
|
||||
logger.level=DEBUG
|
||||
logger.file.enabled=true
|
||||
logger.file=logs/application.log
|
||||
logger.max.size.mb=10
|
||||
@@ -0,0 +1,56 @@
|
||||
import paho.mqtt.client as mqtt
|
||||
import sys
|
||||
import random
|
||||
import time
|
||||
|
||||
# ===== KONFIGURATION =====
|
||||
BROKER = "127.0.0.1"
|
||||
PORT = 1883
|
||||
TOPIC = "PREDICTION"
|
||||
USERNAME = None
|
||||
PASSWORD = None
|
||||
QOS = 0
|
||||
INTERVAL_SECONDS = 10
|
||||
# ==========================
|
||||
|
||||
|
||||
def on_connect(client, userdata, flags, rc):
|
||||
if rc == 0:
|
||||
print("Erfolgreich mit Broker verbunden")
|
||||
else:
|
||||
print(f"Verbindung fehlgeschlagen mit Code {rc}")
|
||||
|
||||
|
||||
def main():
|
||||
client = mqtt.Client()
|
||||
|
||||
if USERNAME and PASSWORD:
|
||||
client.username_pw_set(USERNAME, PASSWORD)
|
||||
|
||||
client.on_connect = on_connect
|
||||
|
||||
try:
|
||||
client.connect(BROKER, PORT, 60)
|
||||
client.loop_start() # Non-blocking loop
|
||||
|
||||
print("Starte kontinuierliches Senden... (STRG+C zum Beenden)")
|
||||
|
||||
while True:
|
||||
message = random.randint(0, 1)
|
||||
client.publish(TOPIC, message, qos=QOS)
|
||||
print(f"Gesendet an '{TOPIC}': {message}")
|
||||
time.sleep(INTERVAL_SECONDS)
|
||||
|
||||
except KeyboardInterrupt:
|
||||
print("\nBeende Publisher...")
|
||||
client.loop_stop()
|
||||
client.disconnect()
|
||||
sys.exit(0)
|
||||
|
||||
except Exception as e:
|
||||
print("Fehler:", e)
|
||||
sys.exit(1)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,49 @@
|
||||
package vassistent.service;
|
||||
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.mockito.Mockito.*;
|
||||
|
||||
class BinaryEventServiceTest {
|
||||
|
||||
private DataPersistenceService persistenceService;
|
||||
private EvaluationService evaluationService;
|
||||
private BinaryEventService binaryEventService;
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
persistenceService = mock(DataPersistenceService.class);
|
||||
evaluationService = mock(EvaluationService.class);
|
||||
binaryEventService = new BinaryEventService(persistenceService, evaluationService);
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
persistenceService = null;
|
||||
evaluationService = null;
|
||||
binaryEventService = null;
|
||||
}
|
||||
|
||||
@Test
|
||||
void handlePayloadValid() {
|
||||
binaryEventService.handlePayload("1");
|
||||
verify(persistenceService).store(anyInt());
|
||||
verify(evaluationService).evaluate();
|
||||
}
|
||||
|
||||
@Test
|
||||
void handlePayloadInvalidValue() {
|
||||
binaryEventService.handlePayload("5");
|
||||
verify(persistenceService, never()).store(anyInt());
|
||||
verify(evaluationService, never()).evaluate();
|
||||
}
|
||||
|
||||
@Test
|
||||
void handlePayloadNonNumeric() {
|
||||
binaryEventService.handlePayload("abc");
|
||||
verify(persistenceService, never()).store(anyInt());
|
||||
verify(evaluationService, never()).evaluate();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package vassistent.service;
|
||||
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
|
||||
import java.io.File;
|
||||
import java.sql.Connection;
|
||||
import java.sql.DriverManager;
|
||||
import java.sql.ResultSet;
|
||||
import java.sql.Statement;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
class DataPersistenceServiceTest {
|
||||
|
||||
private DataPersistenceService service;
|
||||
private static final String TEST_DB = "jdbc:sqlite:data/test_health.db";
|
||||
private static final String DB_FILE = "data/test_health.db";
|
||||
|
||||
@BeforeEach
|
||||
void setup() {
|
||||
service = new DataPersistenceService(TEST_DB);
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
void cleanup() {
|
||||
deleteTestDatabase();
|
||||
}
|
||||
|
||||
|
||||
private void deleteTestDatabase() {
|
||||
File db = new File(DB_FILE);
|
||||
if (db.exists()) {
|
||||
assertTrue(db.delete(), "Testdatenbank konnte nicht gelöscht werden");
|
||||
}
|
||||
}
|
||||
|
||||
@org.junit.jupiter.api.Test
|
||||
void createDatabaseFolder() {
|
||||
File folder = new File("data");
|
||||
|
||||
assertTrue(folder.exists(), "Datenbankordner sollte existieren");
|
||||
}
|
||||
|
||||
@org.junit.jupiter.api.Test
|
||||
void store() throws Exception {
|
||||
service.store(1);
|
||||
|
||||
try (Connection c = DriverManager.getConnection(TEST_DB);
|
||||
Statement s = c.createStatement();
|
||||
ResultSet rs = s.executeQuery(
|
||||
"SELECT COUNT(*) FROM binary_event")) {
|
||||
|
||||
assertTrue(rs.next());
|
||||
int count = rs.getInt(1);
|
||||
|
||||
assertTrue(count > 0,
|
||||
"Es sollte mindestens ein Eintrag existieren");
|
||||
}
|
||||
}
|
||||
|
||||
@org.junit.jupiter.api.Test
|
||||
void getDbUrl() {
|
||||
assertEquals(TEST_DB, service.getDbUrl());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
package vassistent.service;
|
||||
|
||||
import org.junit.jupiter.api.Assertions;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import vassistent.model.AppState;
|
||||
import vassistent.model.ProblemLevel;
|
||||
|
||||
import static org.mockito.Mockito.*;
|
||||
|
||||
class EvaluationServiceTest {
|
||||
|
||||
private StatisticsService statisticsService;
|
||||
private AppState appState;
|
||||
private UnrealWebSocketService unrealService;
|
||||
private EvaluationService evaluationService;
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
statisticsService = mock(StatisticsService.class);
|
||||
appState = new AppState();
|
||||
unrealService = mock(UnrealWebSocketService.class);
|
||||
evaluationService = new EvaluationService(statisticsService, appState, unrealService);
|
||||
}
|
||||
|
||||
@org.junit.jupiter.api.Test
|
||||
void tearDown() {
|
||||
statisticsService = null;
|
||||
appState = null;
|
||||
unrealService = null;
|
||||
evaluationService = null;
|
||||
}
|
||||
|
||||
@org.junit.jupiter.api.Test
|
||||
void testEvaluateDisasterLevel() {
|
||||
when(statisticsService.getRatio(anyInt())).thenReturn(0.95);
|
||||
evaluationService.evaluate();
|
||||
Assertions.assertEquals(ProblemLevel.DISASTER, appState.getProblemLevel());
|
||||
verify(unrealService).speak("DISASTER");
|
||||
}
|
||||
|
||||
@Test
|
||||
void testEvaluateNoChange() {
|
||||
appState.setProblemLevel(ProblemLevel.NONE);
|
||||
when(statisticsService.getRatio(anyInt())).thenReturn(0.1);
|
||||
evaluationService.evaluate();
|
||||
verify(unrealService, never()).speak(anyString());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
package vassistent.service;
|
||||
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.io.File;
|
||||
import java.sql.SQLException;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
class StatisticsServiceTest {
|
||||
|
||||
private String DB_URL;
|
||||
private DataPersistenceService persistenceService;
|
||||
private StatisticsService statisticsService;
|
||||
private File tempDbFile;
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
try {
|
||||
tempDbFile = File.createTempFile("testdb", ".db");
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException("Kann temporäre DB-Datei nicht erstellen", e);
|
||||
}
|
||||
|
||||
DB_URL = "jdbc:sqlite:" + tempDbFile.getAbsolutePath();
|
||||
persistenceService = new DataPersistenceService(DB_URL);
|
||||
statisticsService = new StatisticsService(persistenceService);
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
if (tempDbFile != null && tempDbFile.exists()) {
|
||||
tempDbFile.delete();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void testGetRatioNoData() {
|
||||
double ratio = statisticsService.getRatio(10);
|
||||
assertEquals(0.0, ratio, "Ratio sollte 0 sein, wenn keine Daten vorhanden");
|
||||
}
|
||||
|
||||
@Test
|
||||
void testGetRatioSingleValue() throws SQLException {
|
||||
persistenceService.store(1);
|
||||
|
||||
double ratio = statisticsService.getRatio(1);
|
||||
assertEquals(1.0, ratio, "Ratio sollte 1 sein, wenn nur 1 gespeichert wurde");
|
||||
}
|
||||
|
||||
@Test
|
||||
void testGetRatioMultipleValues() {
|
||||
persistenceService.store(1);
|
||||
persistenceService.store(0);
|
||||
persistenceService.store(1);
|
||||
|
||||
double ratio = statisticsService.getRatio(10);
|
||||
assertEquals((1 + 0 + 1) / 3.0, ratio, 0.0001, "Ratio sollte Durchschnitt der letzten Werte sein");
|
||||
}
|
||||
|
||||
@Test
|
||||
void testGetRatioLimitLastN() {
|
||||
persistenceService.store(1);
|
||||
persistenceService.store(1);
|
||||
persistenceService.store(0);
|
||||
persistenceService.store(0);
|
||||
persistenceService.store(1);
|
||||
|
||||
double ratio = statisticsService.getRatio(3);
|
||||
assertEquals((1 + 0 + 0) / 3.0, ratio, 0.0001, "Ratio sollte Durchschnitt der letzten Werte sein");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user