first commit
This commit is contained in:
206
lib/utils/mqtt_service.dart
Normal file
206
lib/utils/mqtt_service.dart
Normal file
@@ -0,0 +1,206 @@
|
||||
import 'dart:async';
|
||||
import 'dart:convert';
|
||||
import 'package:flutter/foundation.dart';
|
||||
import 'package:mqtt_client/mqtt_client.dart';
|
||||
import 'package:mqtt_client/mqtt_server_client.dart';
|
||||
import 'package:miler/views/helpers/constants/mqtt_constants.dart';
|
||||
import 'package:shared_preferences/shared_preferences.dart';
|
||||
|
||||
class MilerMqttService {
|
||||
static final MilerMqttService _instance = MilerMqttService._internal();
|
||||
|
||||
factory MilerMqttService() {
|
||||
return _instance;
|
||||
}
|
||||
|
||||
MilerMqttService._internal();
|
||||
|
||||
MqttServerClient? _client;
|
||||
bool _isConnected = false;
|
||||
String? _currentRiderId;
|
||||
|
||||
bool get isConnected => _isConnected;
|
||||
|
||||
Future<void> connect() async {
|
||||
if (_isConnected) return;
|
||||
|
||||
final prefs = await SharedPreferences.getInstance();
|
||||
final riderIdLong = prefs.getInt('userId') ?? prefs.getInt('userid');
|
||||
if (riderIdLong == null || riderIdLong <= 0) {
|
||||
debugPrint('[MQTT] Cannot connect: No rider ID found in preferences.');
|
||||
return;
|
||||
}
|
||||
_currentRiderId = riderIdLong.toString();
|
||||
|
||||
const host = MqttConstants.brokerHost;
|
||||
// Append a small random or "bg" suffix if we want to allow simultaneous connections
|
||||
// from UI and Background Isolates.
|
||||
final clientId =
|
||||
'rider_${_currentRiderId}_${DateTime.now().millisecondsSinceEpoch % 1000}';
|
||||
|
||||
_client = MqttServerClient(host, clientId);
|
||||
_client!.port = MqttConstants.brokerPort;
|
||||
_client!.keepAlivePeriod = 30; // 30 seconds
|
||||
_client!.autoReconnect = true;
|
||||
_client!.logging(on: false);
|
||||
|
||||
// --- LAST WILL AND TESTAMENT (LWT) ---
|
||||
// If the app crashes or loses internet unexpectedly, the broker will publish this.
|
||||
final lwtTopic = MqttConstants.topicRiderStatus.replaceAll(
|
||||
'{riderId}',
|
||||
_currentRiderId!,
|
||||
);
|
||||
_client!.onDisconnected = _onDisconnected;
|
||||
_client!.onConnected = _onConnected;
|
||||
_client!.onAutoReconnect = _onAutoReconnect;
|
||||
_client!.onSubscribed = _onSubscribed;
|
||||
|
||||
final connMessage = MqttConnectMessage()
|
||||
.withClientIdentifier(clientId)
|
||||
.authenticateAs(MqttConstants.username, MqttConstants.passwordString)
|
||||
.withWillTopic(lwtTopic)
|
||||
.withWillMessage(MqttConstants.statusOffline)
|
||||
.withWillQos(MqttQos.atLeastOnce)
|
||||
.withWillRetain()
|
||||
.startClean();
|
||||
|
||||
_client!.connectionMessage = connMessage;
|
||||
|
||||
try {
|
||||
debugPrint('[MQTT] Connecting to $host...');
|
||||
await _client!.connect();
|
||||
} catch (e) {
|
||||
debugPrint('[MQTT] Connection failed: $e');
|
||||
_disconnect();
|
||||
}
|
||||
}
|
||||
|
||||
void _onConnected() {
|
||||
_isConnected = true;
|
||||
debugPrint('[MQTT] Connected successfully.');
|
||||
_publishStatus(MqttConstants.statusOnline);
|
||||
}
|
||||
|
||||
void _onDisconnected() {
|
||||
_isConnected = false;
|
||||
debugPrint('[MQTT] Disconnected from broker.');
|
||||
}
|
||||
|
||||
void _onAutoReconnect() {
|
||||
debugPrint('[MQTT] Auto-reconnecting...');
|
||||
}
|
||||
|
||||
void _onSubscribed(String topic) {
|
||||
debugPrint('[MQTT] Subscribed to topic: $topic');
|
||||
}
|
||||
|
||||
void _disconnect() {
|
||||
_client?.disconnect();
|
||||
_onDisconnected();
|
||||
}
|
||||
|
||||
void _publishStatus(String status) {
|
||||
if (!_isConnected || _currentRiderId == null) return;
|
||||
|
||||
final topic = MqttConstants.topicRiderStatus.replaceAll(
|
||||
'{riderId}',
|
||||
_currentRiderId!,
|
||||
);
|
||||
final builder = MqttClientPayloadBuilder();
|
||||
builder.addString(status);
|
||||
|
||||
_client!.publishMessage(
|
||||
topic,
|
||||
MqttQos.atLeastOnce,
|
||||
builder.payload!,
|
||||
retain: true,
|
||||
);
|
||||
debugPrint('[MQTT] Published status: $status to $topic');
|
||||
}
|
||||
|
||||
// --- PUBLIC API ---
|
||||
|
||||
/// Updates the rider's status (Online, Offline, Idle)
|
||||
void updateStatus(String status) {
|
||||
_publishStatus(status);
|
||||
}
|
||||
|
||||
/// Updates the rider's profile (name, context, tokens etc.)
|
||||
void publishProfile(Map<String, dynamic> profileData) {
|
||||
if (!_isConnected || _currentRiderId == null) return;
|
||||
|
||||
final topic = MqttConstants.topicRiderProfile.replaceAll(
|
||||
'{riderId}',
|
||||
_currentRiderId!,
|
||||
);
|
||||
final builder = MqttClientPayloadBuilder();
|
||||
builder.addString(jsonEncode(profileData));
|
||||
|
||||
_client!.publishMessage(
|
||||
topic,
|
||||
MqttQos.atLeastOnce,
|
||||
builder.payload!,
|
||||
retain: true,
|
||||
);
|
||||
debugPrint('[MQTT] Published profile to $topic');
|
||||
}
|
||||
|
||||
/// Publishes location data
|
||||
void publishLocation(Map<String, dynamic> locationData) {
|
||||
if (!_isConnected || _currentRiderId == null) return;
|
||||
|
||||
final topic = MqttConstants.topicRiderLocation.replaceAll(
|
||||
'{riderId}',
|
||||
_currentRiderId!,
|
||||
);
|
||||
final builder = MqttClientPayloadBuilder();
|
||||
builder.addString(jsonEncode(locationData));
|
||||
|
||||
_client!.publishMessage(topic, MqttQos.atMostOnce, builder.payload!);
|
||||
// Avoid excessive logging for high-frequency location updates
|
||||
}
|
||||
|
||||
/// Publishes telemetry data (Battery, GPS Signal, Rider Logs etc.)
|
||||
void publishTelemetry(Map<String, dynamic> telemetryData) {
|
||||
if (!_isConnected || _currentRiderId == null) return;
|
||||
|
||||
final topic = MqttConstants.topicRiderTelemetry.replaceAll(
|
||||
'{riderId}',
|
||||
_currentRiderId!,
|
||||
);
|
||||
final builder = MqttClientPayloadBuilder();
|
||||
builder.addString(jsonEncode(telemetryData));
|
||||
|
||||
_client!.publishMessage(topic, MqttQos.atLeastOnce, builder.payload!);
|
||||
debugPrint('[MQTT] Published telemetry data.');
|
||||
}
|
||||
|
||||
/// Publishes arbitrary logs or events
|
||||
void publishLog(String eventName, Map<String, dynamic> data) {
|
||||
if (!_isConnected || _currentRiderId == null) return;
|
||||
|
||||
final topic =
|
||||
'${MqttConstants.topicRiderLogs.replaceAll('{riderId}', _currentRiderId!)}/$eventName';
|
||||
final builder = MqttClientPayloadBuilder();
|
||||
builder.addString(jsonEncode(data));
|
||||
|
||||
_client!.publishMessage(topic, MqttQos.atLeastOnce, builder.payload!);
|
||||
debugPrint('[MQTT] Published log: $eventName');
|
||||
}
|
||||
|
||||
/// Generic publish to doormile/riders/{id}/{subTopic}
|
||||
void publish(String subTopic, dynamic data) {
|
||||
if (!_isConnected || _currentRiderId == null) return;
|
||||
|
||||
final topic = 'doormile/riders/$_currentRiderId/$subTopic';
|
||||
final builder = MqttClientPayloadBuilder();
|
||||
if (data is String) {
|
||||
builder.addString(data);
|
||||
} else {
|
||||
builder.addString(jsonEncode(data));
|
||||
}
|
||||
|
||||
_client!.publishMessage(topic, MqttQos.atLeastOnce, builder.payload!);
|
||||
debugPrint('[MQTT] Published to $topic');
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user