MMCL/SOURCE/agents/delphi_nilm_agent/__history/uMain.pas.~9~
2026-09-04 11:27:31 +09:00

279 lines
8.7 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; // JSON 파싱, Map 컬렉션, MQTT 연동
type
TfMain = class(TForm)
btnStart: TButton;
btnStop: TButton;
memLog: TMemo;
procedure btnStartClick(Sender: TObject);
procedure btnStopClick(Sender: TObject);
private
FClient: TMQTTClient;
FTopicSensorMap: TDictionary<string, Integer>;
procedure SetupDatabase;
procedure LogMsg(const AMsg: string);
procedure InsertNilmData(SensorNo: Integer; PayloadJSON: string);
// 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.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;
procedure TfMain.InsertNilmData(SensorNo: Integer; PayloadJSON: string);
var
Qry: TFDQuery;
JSONObj: TJSONObject;
v1_pwr, v1_cur, v1_volt: Double;
v2_pwr, v2_cur, v2_volt: Double;
v3_pwr, v3_cur, v3_volt: Double;
begin
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 파싱 (System.JSON 활용)
JSONObj := TJSONObject.ParseJSONValue(PayloadJSON) as TJSONObject;
if Assigned(JSONObj) then
begin
try
// CH1 (L1 / A)
JSONObj.TryGetValue<Double>('wA', v1_pwr);
JSONObj.TryGetValue<Double>('currentA', v1_cur);
JSONObj.TryGetValue<Double>('voltageA', v1_volt);
// CH2 (L2 / B)
JSONObj.TryGetValue<Double>('wB', v2_pwr);
JSONObj.TryGetValue<Double>('currentB', v2_cur);
JSONObj.TryGetValue<Double>('voltageB', v2_volt);
// CH3 (L3 / C)
JSONObj.TryGetValue<Double>('wC', v3_pwr);
JSONObj.TryGetValue<Double>('currentC', v3_cur);
JSONObj.TryGetValue<Double>('voltageC', v3_volt);
finally
JSONObj.Free;
end;
end;
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]));
except
on E: Exception do
LogMsg('DB 처리 오류: ' + E.Message);
end;
Qry.Free;
end;
procedure TfMain.OnMqttMessageReceived(const Topic, Payload: string);
var
SensorNo: Integer;
begin
LogMsg('MQTT 수신 Topic: ' + Topic);
if Assigned(FTopicSensorMap) and FTopicSensorMap.TryGetValue(Topic, SensorNo) then
begin
InsertNilmData(SensorNo, Payload);
end
else
begin
LogMsg('DB에 맵핑되지 않은 토픽의 메시지 무시: ' + Topic);
end;
end;
procedure TfMain.OnMqttStatusChanged(AConnected: Boolean; const AErrorMsg: string);
begin
if AConnected then
LogMsg('MQTT 브로커 접속 성공 (구독 정보 서버 전송 완료)')
else
LogMsg('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', '');
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;
// 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 DM.fdConnEtc.Connected then
begin
DM.fdConnEtc.Connected := False;
LogMsg('DB 연결이 안전하게 해제되었습니다.');
end;
end;
end.