446 lines
14 KiB
Plaintext
446 lines
14 KiB
Plaintext
unit uMain;
|
|
|
|
interface
|
|
|
|
uses
|
|
Winapi.Windows, Winapi.Messages, System.SysUtils, System.Variants, System.Classes, Vcl.Graphics,
|
|
Vcl.Controls, Vcl.Forms, Vcl.Dialogs, Vcl.StdCtrls,
|
|
System.JSON, System.Generics.Collections, UMQTTClient, Vcl.ComCtrls, System.Net.HttpClient, System.Net.URLClient;
|
|
|
|
type
|
|
TNilmValues = record
|
|
v1_volt, v1_cur, v1_pwr: Double;
|
|
v2_volt, v2_cur, v2_pwr: Double;
|
|
v3_volt, v3_cur, v3_pwr: Double;
|
|
end;
|
|
|
|
TfMain = class(TForm)
|
|
btnStart: TButton;
|
|
btnStop: TButton;
|
|
PageControl1: TPageControl;
|
|
TabSheet1: TTabSheet;
|
|
TabSheet2: TTabSheet;
|
|
memLog: TMemo;
|
|
memLogMqtt: TMemo;
|
|
procedure btnStartClick(Sender: TObject);
|
|
procedure btnStopClick(Sender: TObject);
|
|
private
|
|
FClient: TMQTTClient;
|
|
FTopicSensorMap: TDictionary<string, Integer>;
|
|
FLastValues: TDictionary<Integer, TNilmValues>;
|
|
FAnalyzedPwr: TDictionary<Integer, Double>;
|
|
FBackendUrl: string;
|
|
|
|
procedure SetupDatabase;
|
|
procedure LogMsg(const AMsg: string);
|
|
procedure LogMqttMsg(const AMsg: string);
|
|
procedure AnalyzePowerAsync(SensorNo: Integer; CurrentPwr: Double);
|
|
function InsertNilmData(SensorNo: Integer; PayloadJSON: string): Boolean;
|
|
|
|
// MQTT 수신 이벤트 콜백 시그니처 (컴포넌트에 맞게 수정하세요)
|
|
procedure OnMqttMessageReceived(const Topic, Payload: string);
|
|
procedure OnMqttStatusChanged(AConnected: Boolean; const AErrorMsg: string);
|
|
public
|
|
{ Public declarations }
|
|
end;
|
|
|
|
var
|
|
fMain: TfMain;
|
|
|
|
implementation
|
|
|
|
{$R *.dfm}
|
|
|
|
uses
|
|
System.IniFiles,
|
|
U_DM, FireDAC.Comp.Client;
|
|
|
|
procedure TfMain.LogMsg(const AMsg: string);
|
|
begin
|
|
// 멀티스레딩 환경(MQTT 콜백 등)에서 VCL 메모 컨트롤 접근 시 에러 방지
|
|
TThread.Queue(nil,
|
|
procedure
|
|
begin
|
|
memLog.Lines.Add(FormatDateTime('yyyy-mm-dd hh:nn:ss', Now) + ' - ' + AMsg);
|
|
end);
|
|
end;
|
|
|
|
procedure TfMain.LogMqttMsg(const AMsg: string);
|
|
begin
|
|
TThread.Queue(nil,
|
|
procedure
|
|
begin
|
|
memLogMqtt.Lines.Add(FormatDateTime('yyyy-mm-dd hh:nn:ss', Now) + ' - ' + AMsg);
|
|
end);
|
|
end;
|
|
|
|
procedure TfMain.AnalyzePowerAsync(SensorNo: Integer; CurrentPwr: Double);
|
|
begin
|
|
TThread.CreateAnonymousThread(
|
|
procedure
|
|
var
|
|
HttpClient: TNetHTTPClient;
|
|
JsonBody: TJSONObject;
|
|
Resp: IHTTPResponse;
|
|
SS: TStringStream;
|
|
ParsedResp: TJSONObject;
|
|
begin
|
|
HttpClient := TNetHTTPClient.Create(nil);
|
|
JsonBody := TJSONObject.Create;
|
|
SS := TStringStream.Create('', TEncoding.UTF8);
|
|
try
|
|
JsonBody.AddPair('sensor_no', TJSONNumber.Create(SensorNo));
|
|
JsonBody.AddPair('current_power', TJSONNumber.Create(CurrentPwr));
|
|
SS.WriteString(JsonBody.ToString);
|
|
SS.Position := 0;
|
|
|
|
try
|
|
TThread.Queue(nil, procedure begin LogMqttMsg(Format('[AI 상태 분석 요청] Sensor: %d, Power: %.1fW', [SensorNo, CurrentPwr])); end);
|
|
Resp := HttpClient.Post(FBackendUrl + '/api/analyze_power', SS, nil,
|
|
[TNameValuePair.Create('Content-Type', 'application/json')]);
|
|
|
|
if Resp.StatusCode = 200 then
|
|
begin
|
|
ParsedResp := TJSONObject.ParseJSONValue(Resp.ContentAsString(TEncoding.UTF8)) as TJSONObject;
|
|
if Assigned(ParsedResp) then
|
|
try
|
|
TThread.Queue(nil, procedure begin LogMqttMsg(ParsedResp.GetValue('reply').Value); end);
|
|
finally
|
|
ParsedResp.Free;
|
|
end;
|
|
end
|
|
else
|
|
TThread.Queue(nil, procedure begin LogMqttMsg('AI 전력 분석 호출 실패 (HTTP ' + IntToStr(Resp.StatusCode) + ')'); end);
|
|
except
|
|
on E: Exception do
|
|
TThread.Queue(nil, procedure begin LogMqttMsg('AI 전력 분석 연결 오류: ' + E.Message); end);
|
|
end;
|
|
finally
|
|
SS.Free;
|
|
JsonBody.Free;
|
|
HttpClient.Free;
|
|
end;
|
|
end
|
|
).Start;
|
|
end;
|
|
|
|
procedure TfMain.SetupDatabase;
|
|
var
|
|
Ini: TIniFile;
|
|
DBHost, DBUser, DBPass, DBName: string;
|
|
DBPort: Integer;
|
|
begin
|
|
if DM.fdConnEtc.Connected then Exit;
|
|
|
|
Ini := TIniFile.Create(ExtractFilePath(ParamStr(0)) + 'settings.ini');
|
|
try
|
|
DBHost := Ini.ReadString('DB', 'Host', '127.0.0.1');
|
|
DBPort := Ini.ReadInteger('DB', 'Port', 3306);
|
|
DBUser := Ini.ReadString('DB', 'User', 'root');
|
|
DBPass := Ini.ReadString('DB', 'Password', 'password');
|
|
DBName := Ini.ReadString('DB', 'Database', 'mmcl_db');
|
|
finally
|
|
Ini.Free;
|
|
end;
|
|
|
|
// DM.FDPhysPgDriverLink.VendorLib := '.\libpq-10.dll';
|
|
|
|
// U_DM 컴포넌트를 설계서 상의 MariaDB/MySQL 구동 형식으로 연결
|
|
DM.fdConnEtc.Params.Clear;
|
|
DM.fdConnEtc.Params.DriverID := 'MySQL';
|
|
DM.fdConnEtc.Params.Database := DBName;
|
|
DM.fdConnEtc.Params.UserName := DBUser;
|
|
DM.fdConnEtc.Params.Password := DBPass;
|
|
DM.fdConnEtc.Params.Add('Server=' + DBHost);
|
|
DM.fdConnEtc.Params.Add('Port=' + IntToStr(DBPort));
|
|
|
|
try
|
|
DM.fdConnEtc.Connected := True;
|
|
LogMsg('MariaDB 연결 성공 (Host: ' + DBHost + ') - U_DM 연동 완료');
|
|
except
|
|
on E: Exception do
|
|
LogMsg('DB 연결 오류: ' + E.Message);
|
|
end;
|
|
end;
|
|
|
|
function TfMain.InsertNilmData(SensorNo: Integer; PayloadJSON: string): Boolean;
|
|
var
|
|
Qry: TFDQuery;
|
|
JSONRoot: TJSONObject;
|
|
SlaveDataArr: TJSONArray;
|
|
JSONObj: TJSONObject;
|
|
v1_pwr, v1_cur, v1_volt: Double;
|
|
v2_pwr, v2_cur, v2_volt: Double;
|
|
v3_pwr, v3_cur, v3_volt: Double;
|
|
CurrentValues, LastValues: TNilmValues;
|
|
Changed: Boolean;
|
|
begin
|
|
Result := False;
|
|
if not DM.fdConnEtc.Connected then Exit;
|
|
|
|
// 기본값 초기화
|
|
v1_pwr:=0; v1_cur:=0; v1_volt:=0;
|
|
v2_pwr:=0; v2_cur:=0; v2_volt:=0;
|
|
v3_pwr:=0; v3_cur:=0; v3_volt:=0;
|
|
|
|
// JSON 파싱 (새로운 구조: {"slave_data": [{...}]})
|
|
try
|
|
JSONRoot := TJSONObject.ParseJSONValue(PayloadJSON) as TJSONObject;
|
|
if Assigned(JSONRoot) then
|
|
begin
|
|
try
|
|
if JSONRoot.TryGetValue<TJSONArray>('slave_data', SlaveDataArr) and (SlaveDataArr.Count > 0) then
|
|
begin
|
|
JSONObj := SlaveDataArr.Items[0] as TJSONObject;
|
|
if Assigned(JSONObj) then
|
|
begin
|
|
// CH1 (L1 / A)
|
|
JSONObj.TryGetValue<Double>('Va_A', v1_pwr);
|
|
JSONObj.TryGetValue<Double>('currentA', v1_cur);
|
|
JSONObj.TryGetValue<Double>('voltageA', v1_volt);
|
|
// CH2 (L2 / B)
|
|
JSONObj.TryGetValue<Double>('Va_B', v2_pwr);
|
|
JSONObj.TryGetValue<Double>('currentB', v2_cur);
|
|
JSONObj.TryGetValue<Double>('voltageB', v2_volt);
|
|
// CH3 (L3 / C) - 3상일 경우 존재함
|
|
JSONObj.TryGetValue<Double>('Va_C', v3_pwr);
|
|
JSONObj.TryGetValue<Double>('currentC', v3_cur);
|
|
JSONObj.TryGetValue<Double>('voltageC', v3_volt);
|
|
end;
|
|
end;
|
|
finally
|
|
JSONRoot.Free;
|
|
end;
|
|
end;
|
|
except
|
|
on E: Exception do
|
|
begin
|
|
LogMqttMsg('JSON 파싱 오류: ' + E.Message);
|
|
Exit;
|
|
end;
|
|
end;
|
|
|
|
CurrentValues.v1_pwr := v1_pwr; CurrentValues.v1_cur := v1_cur; CurrentValues.v1_volt := v1_volt;
|
|
CurrentValues.v2_pwr := v2_pwr; CurrentValues.v2_cur := v2_cur; CurrentValues.v2_volt := v2_volt;
|
|
CurrentValues.v3_pwr := v3_pwr; CurrentValues.v3_cur := v3_cur; CurrentValues.v3_volt := v3_volt;
|
|
|
|
Changed := True;
|
|
if Assigned(FLastValues) and FLastValues.TryGetValue(SensorNo, LastValues) then
|
|
begin
|
|
if (CurrentValues.v1_pwr = LastValues.v1_pwr) and (CurrentValues.v1_cur = LastValues.v1_cur) and (CurrentValues.v1_volt = LastValues.v1_volt) and
|
|
(CurrentValues.v2_pwr = LastValues.v2_pwr) and (CurrentValues.v2_cur = LastValues.v2_cur) and (CurrentValues.v2_volt = LastValues.v2_volt) and
|
|
(CurrentValues.v3_pwr = LastValues.v3_pwr) and (CurrentValues.v3_cur = LastValues.v3_cur) and (CurrentValues.v3_volt = LastValues.v3_volt) then
|
|
begin
|
|
Changed := False;
|
|
end;
|
|
end;
|
|
|
|
if not Changed then
|
|
begin
|
|
// 값의 변경이 없으므로 DB 로직 수행 안함
|
|
Result := False;
|
|
Exit;
|
|
end;
|
|
|
|
Result := True;
|
|
|
|
if Assigned(FLastValues) then
|
|
FLastValues.AddOrSetValue(SensorNo, CurrentValues);
|
|
|
|
Qry := TFDQuery.Create(nil);
|
|
try
|
|
Qry.Connection := DM.fdConnEtc;
|
|
|
|
// 1. sensor_history_log (시계열 이력) 적재 - CH1, CH2, CH3 모두 저장
|
|
Qry.SQL.Text :=
|
|
'INSERT INTO sensor_history_log (sensor_no, log_time, ' +
|
|
' value_ch1_pwr, value_ch1_current, value_ch1_volt, ' +
|
|
' value_ch2_pwr, value_ch2_current, value_ch2_volt, ' +
|
|
' value_ch3_pwr, value_ch3_current, value_ch3_volt) ' +
|
|
'VALUES (:SNo, NOW(), :P1, :C1, :V1, :P2, :C2, :V2, :P3, :C3, :V3)';
|
|
|
|
Qry.ParamByName('SNo').AsInteger := SensorNo;
|
|
Qry.ParamByName('P1').AsFloat := v1_pwr; Qry.ParamByName('C1').AsFloat := v1_cur; Qry.ParamByName('V1').AsFloat := v1_volt;
|
|
Qry.ParamByName('P2').AsFloat := v2_pwr; Qry.ParamByName('C2').AsFloat := v2_cur; Qry.ParamByName('V2').AsFloat := v2_volt;
|
|
Qry.ParamByName('P3').AsFloat := v3_pwr; Qry.ParamByName('C3').AsFloat := v3_cur; Qry.ParamByName('V3').AsFloat := v3_volt;
|
|
Qry.ExecSQL;
|
|
|
|
// 2. sensor_info (현재 상태) 업데이트
|
|
Qry.SQL.Text :=
|
|
'UPDATE sensor_info SET update_time = NOW(), ' +
|
|
' value_ch1_pwr = :P1, value_ch1_current = :C1, value_ch1_volt = :V1, ' +
|
|
' value_ch2_pwr = :P2, value_ch2_current = :C2, value_ch2_volt = :V2, ' +
|
|
' value_ch3_pwr = :P3, value_ch3_current = :C3, value_ch3_volt = :V3 ' +
|
|
'WHERE sensor_no = :SNo AND sensor_typeid = 1';
|
|
|
|
Qry.ParamByName('SNo').AsInteger := SensorNo;
|
|
Qry.ParamByName('P1').AsFloat := v1_pwr; Qry.ParamByName('C1').AsFloat := v1_cur; Qry.ParamByName('V1').AsFloat := v1_volt;
|
|
Qry.ParamByName('P2').AsFloat := v2_pwr; Qry.ParamByName('C2').AsFloat := v2_cur; Qry.ParamByName('V2').AsFloat := v2_volt;
|
|
Qry.ParamByName('P3').AsFloat := v3_pwr; Qry.ParamByName('C3').AsFloat := v3_cur; Qry.ParamByName('V3').AsFloat := v3_volt;
|
|
Qry.ExecSQL;
|
|
|
|
LogMsg(Format('DB Insert/Update 완료 [Sensor: %d]', [SensorNo]));
|
|
|
|
// 유의미한 전력 변동 감지 (예: 1000W 이상 변동 시, 혹은 20% 이상 변동 시 등)
|
|
// 여기서는 절대값 차이가 500W 이상이거나, 처음 보내는 경우 백엔드 분석 호출
|
|
var LastPwr: Double := 0;
|
|
var NeedsAnalysis: Boolean := False;
|
|
|
|
if not Assigned(FAnalyzedPwr) then
|
|
FAnalyzedPwr := TDictionary<Integer, Double>.Create;
|
|
|
|
if FAnalyzedPwr.TryGetValue(SensorNo, LastPwr) then
|
|
begin
|
|
if Abs(CurrentValues.v1_pwr - LastPwr) >= 500.0 then
|
|
NeedsAnalysis := True;
|
|
end
|
|
else
|
|
NeedsAnalysis := True;
|
|
|
|
if NeedsAnalysis then
|
|
begin
|
|
FAnalyzedPwr.AddOrSetValue(SensorNo, CurrentValues.v1_pwr);
|
|
AnalyzePowerAsync(SensorNo, CurrentValues.v1_pwr);
|
|
end;
|
|
|
|
except
|
|
on E: Exception do
|
|
LogMsg('DB 처리 오류: ' + E.Message);
|
|
end;
|
|
Qry.Free;
|
|
end;
|
|
|
|
procedure TfMain.OnMqttMessageReceived(const Topic, Payload: string);
|
|
var
|
|
SensorNo: Integer;
|
|
begin
|
|
if Assigned(FTopicSensorMap) and FTopicSensorMap.TryGetValue(Topic, SensorNo) then
|
|
begin
|
|
if InsertNilmData(SensorNo, Payload) then
|
|
LogMqttMsg('MQTT 수신 [' + Topic + '] 메세지: ' + Payload);
|
|
end
|
|
else
|
|
begin
|
|
// 맵핑 안 된 토픽은 로깅하지 않거나 필요 시 주석 해제
|
|
// LogMqttMsg('DB에 맵핑되지 않은 토픽의 메시지 무시: ' + Topic);
|
|
end;
|
|
end;
|
|
|
|
procedure TfMain.OnMqttStatusChanged(AConnected: Boolean; const AErrorMsg: string);
|
|
begin
|
|
if AConnected then
|
|
LogMqttMsg('MQTT 브로커 접속 성공 (구독 정보 서버 전송 완료)')
|
|
else
|
|
LogMqttMsg('MQTT 브로커 연결 끊김: ' + AErrorMsg);
|
|
end;
|
|
|
|
procedure TfMain.btnStartClick(Sender: TObject);
|
|
var
|
|
Ini: TIniFile;
|
|
MqttHost, MqttUser, MqttPass: string;
|
|
MqttPort: Integer;
|
|
begin
|
|
LogMsg('NILM 에이전트 구동을 시작합니다...');
|
|
SetupDatabase;
|
|
|
|
Ini := TIniFile.Create(ExtractFilePath(ParamStr(0)) + 'settings.ini');
|
|
try
|
|
MqttHost := Ini.ReadString('MQTT', 'Host', '127.0.0.1');
|
|
MqttPort := Ini.ReadInteger('MQTT', 'Port', 1883);
|
|
MqttUser := Ini.ReadString('MQTT', 'User', '');
|
|
MqttPass := Ini.ReadString('MQTT', 'Password', '');
|
|
FBackendUrl := Ini.ReadString('API', 'BackendUrl', 'http://127.0.0.1:8000');
|
|
finally
|
|
Ini.Free;
|
|
end;
|
|
|
|
LogMsg(Format('MQTT 연결 정보 로드 (Host: %s, Port: %d, User: %s)', [MqttHost, MqttPort, MqttUser]));
|
|
|
|
if FTopicSensorMap = nil then
|
|
FTopicSensorMap := TDictionary<string, Integer>.Create
|
|
else
|
|
FTopicSensorMap.Clear;
|
|
|
|
if FLastValues = nil then
|
|
FLastValues := TDictionary<Integer, TNilmValues>.Create
|
|
else
|
|
FLastValues.Clear;
|
|
|
|
if FAnalyzedPwr = nil then
|
|
FAnalyzedPwr := TDictionary<Integer, Double>.Create
|
|
else
|
|
FAnalyzedPwr.Clear;
|
|
|
|
// DB에서 토픽 읽어오기
|
|
if DM.fdConnEtc.Connected then
|
|
begin
|
|
with TFDQuery.Create(nil) do
|
|
try
|
|
Connection := DM.fdConnEtc;
|
|
SQL.Text := 'SELECT sensor_no, mqtt_topic FROM sensor_info WHERE sensor_typeid = 1 AND mqtt_topic IS NOT NULL AND mqtt_topic <> ''''';
|
|
Open;
|
|
while not Eof do
|
|
begin
|
|
FTopicSensorMap.AddOrSetValue(FieldByName('mqtt_topic').AsString, FieldByName('sensor_no').AsInteger);
|
|
Next;
|
|
end;
|
|
LogMsg(Format('DB에서 %d개의 토픽 매핑 정보를 로드했습니다.', [FTopicSensorMap.Count]));
|
|
finally
|
|
Free;
|
|
end;
|
|
end;
|
|
|
|
if Assigned(FClient) then
|
|
begin
|
|
FClient.Disconnect;
|
|
FClient.Free;
|
|
end;
|
|
|
|
FClient := TMQTTClient.Create(MqttHost, MqttPort, '', MqttUser, MqttPass);
|
|
FClient.OnMessage := OnMqttMessageReceived;
|
|
FClient.OnStatus := OnMqttStatusChanged;
|
|
|
|
// 토픽 구독 예약 (Connect 호출 전에 등록해두면 스레드 접속 직후 일괄 구독됨)
|
|
for var Topic in FTopicSensorMap.Keys do
|
|
begin
|
|
FClient.Subscribe(Topic);
|
|
LogMsg('구독 토픽 추가: ' + Topic);
|
|
end;
|
|
|
|
FClient.Connect;
|
|
LogMsg('MQTT 브로커에 연결 요청 완료 (접속 대기중...)');
|
|
end;
|
|
|
|
procedure TfMain.btnStopClick(Sender: TObject);
|
|
begin
|
|
LogMsg('에이전트 중지 요청...');
|
|
|
|
if Assigned(FClient) then
|
|
begin
|
|
FClient.Disconnect;
|
|
FClient.Free;
|
|
FClient := nil;
|
|
LogMsg('MQTT 연결이 안전하게 해제되었습니다.');
|
|
end;
|
|
|
|
if Assigned(FLastValues) then
|
|
begin
|
|
FLastValues.Free;
|
|
FLastValues := nil;
|
|
end;
|
|
|
|
if Assigned(FAnalyzedPwr) then
|
|
begin
|
|
FAnalyzedPwr.Free;
|
|
FAnalyzedPwr := nil;
|
|
end;
|
|
|
|
if DM.fdConnEtc.Connected then
|
|
begin
|
|
DM.fdConnEtc.Connected := False;
|
|
LogMsg('DB 연결이 안전하게 해제되었습니다.');
|
|
end;
|
|
end;
|
|
|
|
end.
|