1"use strict";(self.webpackChunk_N_E=self.webpackChunk_N_E||[]).push([[2129],{65520:(e,t,i)=>{i.d(t,{a:()=>s});class s{getCapabilities(){return{supportsRealTimeStatus:!0,supportsTelemetry:!0,supportsFirmwareManagement:!1,supportsRemoteCommands:!1,supportsBidirectionalSync:!1}}}},92129:(e,t,i)=>{i.d(t,{MqttIntegrationProvider:()=>a});var s=i(41303),n=i(65520),r=i(98409);class a extends n.a{constructor(e){super(),this.providerType="mqtt",this.providerName="MQTT Broker",this.client=null,this.deviceCache=new Map;let t=e.endpoint||e.credentials?.brokerUrl||"",i=e.credentials?.port||1883,s=e.apiKey||e.credentials?.username,n=e.credentials?.password,a=e.credentials?.useTls||!1,o=e.credentials?.clientId||e.credentials?.client_id||`netneural-${Date.now()}`,c=e.credentials?.topicPrefix||"devices/",l=e.credentials?.integrationId||e.projectId||"mqtt",d=e.credentials?.organizationId||"";this.providerId=l,this.integrationId=l,this.organizationId=d,this.activityLogger=new r.m,this.config={brokerUrl:t,port:i,username:s,password:n,useTls:a,clientId:o,topicPrefix:c},this.topicPrefix=c}async testConnection(){return this.organizationId&&this.integrationId?this.activityLogger.withLog({organizationId:this.organizationId,integrationId:this.integrationId,direction:"outgoing",activityType:"test_connection",endpoint:this.config.brokerUrl},async()=>this._testConnectionInternal()):this._testConnectionInternal()}async _testConnectionInternal(){return new Promise(e=>{let t=s.A.connect(this.config.brokerUrl,{username:this.config.username,password:this.config.password,clientId:this.config.clientId||`test_${Date.now()}`,connectTimeout:5e3});t.on("connect",()=>{t.end(),e({success:!0,message:`Successfully connected to MQTT broker: ${this.config.brokerUrl}`})}),t.on("error",i=>{t.end(),e({success:!1,message:i.message||"Failed to connect to MQTT broker"})}),setTimeout(()=>{t.connected||(t.end(),e({success:!1,message:"Connection timeout"}))},5e3)})}async listDevices(e){await this.ensureConnected();let t=Array.from(this.deviceCache.values()).map(e=>e.device),i=e?.page||0,s=e?.limit||100,n=e?.offset||i*s;return{devices:t.slice(n,n+s),total:t.length,page:i,limit:s}}async getDevice(e){await this.ensureConnected();let t=this.deviceCache.get(e);if(!t)throw Error(`Device ${e} not found in cache`);return t.device}async getDeviceStatus(e){await this.ensureConnected();let t=this.deviceCache.get(e);if(!t)throw Error(`Device ${e} not found in cache`);return t.status}async updateDevice(e,t){if(await this.ensureConnected(),!this.client)throw Error("MQTT client not connected");let i=this.deviceCache.get(e);i&&(i.device={...i.device,name:t.name||i.device.name,tags:t.tags||i.device.tags,metadata:{...i.device.metadata,...t.metadata},updatedAt:new Date});let s=this.getDeviceTopic(e,"metadata");return await this.publishAsync(s,JSON.stringify(t)),this.getDevice(e)}async queryTelemetry(){try{let{createClient:e}=await Promise.resolve().then(i.bind(i,17108)),t=e("https://bldojxpockljyivldxwf.supabase.co","eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJpc3MiOiJzdXBhYmFzZSIsInJlZiI6ImJsZG9qeHBvY2tsanlpdmxkeHdmIiwicm9sZSI6ImFub24iLCJpYXQiOjE3NTUwMjY5NTUsImV4cCI6MjA3MDYwMjk1NX0.qkvYx-8ucC5BsqzLcXxIW9TQqc94_dFbGYz5rVSwyRQ"),{data:s,error:n}=await t.from("mqtt_messages").select("*").eq("integration_id",this.integrationId).eq("organization_id",this.organizationId).order("received_at",{ascending:!1}).limit(100);if(n)throw console.error("Error querying MQTT messages:",n),Error(`Failed to query telemetry: ${n.message}`);if(!s||0===s.length)return[];return s.filter(e=>{let t=e.topic;return t.includes("/telemetry")||t.includes("/data")||t.includes("/sensor")}).map(e=>{let t,i=e.payload,s=e.topic,n=s.match(/\/devices\/([^/]+)\//),r=n&&n[1]?n[1]:"unknown";t=new Date(i.timestamp?i.timestamp:i.ts?i.ts:e.received_at);let a={};for(let[e,t]of Object.entries(i))["timestamp","ts","device_id","deviceId"].includes(e)||(a[e]=t);return{deviceId:r,timestamp:t,metrics:a,path:s}})}catch(e){return console.error("Error in queryTelemetry:",e),[]}}getCapabilities(){return{supportsRealTimeStatus:!0,supportsTelemetry:!0,supportsFirmwareManagement:!1,supportsRemoteCommands:!0,supportsBidirectionalSync:!0}}async ensureConnected(){if(!this.client||!this.client.connected)return new Promise((e,t)=>
1{this.client=s.A.connect(this.config.brokerUrl,{username:this.config.username,password:this.config.password,clientId:this.config.clientId||`provider_${this.providerId}`,clean:!1}),this.client.on("connect",()=>{let i=[`${this.topicPrefix}+/status`,`${this.topicPrefix}+/telemetry`,`${this.topicPrefix}+/lwt`];this.client.subscribe(i,i=>{i?t(i):e()})}),this.client.on("message",(e,t)=>{this.handleMessage(e,t)}),this.client.on("error",e=>{t(e)})})}handleMessage(e,t){try{let i=JSON.parse(t.toString()),s=this.extractDeviceIdFromTopic(e);if(!s)return;let n=this.deviceCache.get(s);e.includes("/status")?this.handleStatusMessage(s,i,n):e.includes("/telemetry")?this.handleTelemetryMessage(s,i,n):e.includes("/lwt")&&this.handleLwtMessage(s,i,n)}catch{}}handleStatusMessage(e,t,i){let s=new Date,n="online"===t.status?"online":"offline"===t.status?"offline":"unknown";if(i)i.device.status=n,i.device.lastSeen=s,i.device.updatedAt=s,i.status.connectionState=n,i.status.lastActivity=s,i.lastUpdate=s;else{let i={id:e,name:t.metadata?.name||e,externalId:e,status:n,metadata:t.metadata,lastSeen:s,createdAt:s,updatedAt:s};this.deviceCache.set(e,{device:i,status:{connectionState:n,lastActivity:s,telemetry:{}},lastUpdate:s})}}handleTelemetryMessage(e,t,i){i&&t.telemetry&&(i.status.telemetry={...i.status.telemetry,...t.telemetry},i.lastUpdate=new Date)}handleLwtMessage(e,t,i){i&&(i.device.status="offline",i.status.connectionState="offline",i.lastUpdate=new Date)}extractDeviceIdFromTopic(e){let t=e.match(/devices\/([^/]+)\//);return t?t[1]??null:null}getDeviceTopic(e,t){return`${this.topicPrefix}${e}/${t}`}async publishAsync(e,t){return new Promise((i,s)=>{if(!this.client)return void s(Error("MQTT client not connected"));this.client.publish(e,t,{},e=>{e?s(e):i()})})}async disconnect(){this.client&&(await new Promise(e=>{this.client.end(!1,{},()=>e())}),this.client=null)}}},98409:(e,t,i)=>{i.d(t,{m:()=>n});var s=i(96648);class n{async start(e){try{let{data:t,error:i}=await this.supabase.from("integration_activity_log").insert({organization_id:e.organizationId,integration_id:e.integrationId,direction:e.direction,activity_type:e.activityType,method:e.method,endpoint:e.endpoint,request_body:e.requestBody,status:"started",metadata:e.metadata||{}}).select("id").single();if(i)return console.error("[ActivityLogger] Failed to log activity start:",i),null;return t.id}catch(e){return console.error("[ActivityLogger] Exception logging activity start:",e),null}}async complete(e,t){if(e)try{let{error:i}=await this.supabase.from("integration_activity_log").update({status:t.status,response_status:t.responseStatus,response_body:t.responseBody,response_time_ms:t.responseTimeMs,error_message:t.errorMessage,error_code:t.errorCode,completed_at:new Date().toISOString()}).eq("id",e);i&&console.error("[ActivityLogger] Failed to complete activity log:",i)}catch(e){console.error("[ActivityLogger] Exception completing activity log:",e)}}async withLog(e,t){let i=Date.now(),s=await this.start(e);try{let e=await t(),n=Date.now()-i;return await this.complete(s,{status:"success",responseTimeMs:n,responseBody:"object"==typeof e?e:{value:e}}),e}catch(t){let e=Date.now()-i;throw await this.complete(s,{status:"failed",responseTimeMs:e,errorMessage:t instanceof Error?t.message:"Unknown error",errorCode:t.code}),t}}async getRecentActivity(e){let t=arguments.length>1&&void 0!==arguments[1]?arguments[1]:20,{data:i,error:s}=await this.supabase.from("integration_activity_log").select("*").eq("integration_id",e).order("created_at",{ascending:!1}).limit(t);return s?(console.error("[ActivityLogger] Failed to fetch recent activity:",s),[]):i||[]}async getFailedActivity(e){let t=arguments.length>1&&void 0!==arguments[1]?arguments[1]:10,{data:i,error:s}=await this.supabase.from("integration_activity_log").select("*").eq("integration_id",e).in("status",["failed","error","timeout"]).order("created_at",{ascending:!1}).limit(t);return s?(console.error("[ActivityLogger] Failed to fetch failed activity:",s),[]):i||[]}async getActivityStats(e,t){let i=this.supabase.from("integration_activity_log").select("status, response_time_ms").eq("integration_id",e);t&&(i=i.gte("created_at",t.toISOString()));let{data:s,error:n}=await i;
1if(n||!s)return{total:0,success:0,failed:0,avgResponseTime:null};let r=s.length,a=s.filter(e=>"success"===e.status).length,o=s.filter(e=>["failed","error","timeout"].includes(e.status)).length,c=s.filter(e=>null!==e.response_time_ms).map(e=>e.response_time_ms);return{total:r,success:a,failed:o,avgResponseTime:c.length>0?c.reduce((e,t)=>e+t,0)/c.length:null}}constructor(){this.supabase=(0,s.createClient)()}}}}]);
Line numbers count LF bytes from the start of the resource, as the search results do. Vendor segments are library code the classifier recognised; they are stored but not indexed. Bytes are shown as Latin1 characters, one per byte.