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; FLastValues: TDictionary; FAnalyzedPwr: TDictionary; 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('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('Va_A', v1_pwr); JSONObj.TryGetValue('currentA', v1_cur); JSONObj.TryGetValue('voltageA', v1_volt); // CH2 (L2 / B) JSONObj.TryGetValue('Va_B', v2_pwr); JSONObj.TryGetValue('currentB', v2_cur); JSONObj.TryGetValue('voltageB', v2_volt); // CH3 (L3 / C) - 3상일 경우 존재함 JSONObj.TryGetValue('Va_C', v3_pwr); JSONObj.TryGetValue('currentC', v3_cur); JSONObj.TryGetValue('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.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.Create else FTopicSensorMap.Clear; if FLastValues = nil then FLastValues := TDictionary.Create else FLastValues.Clear; if FAnalyzedPwr = nil then FAnalyzedPwr := TDictionary.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.