feat: add mqtt config model and parser for all Go MQTTOpts fields
This commit is contained in:
@@ -245,6 +245,12 @@ public static class ConfigProcessor
|
||||
opts.ReconnectErrorReports = ToInt(value);
|
||||
break;
|
||||
|
||||
// MQTT
|
||||
case "mqtt":
|
||||
if (value is Dictionary<string, object?> mqttDict)
|
||||
ParseMqtt(mqttDict, opts, errors);
|
||||
break;
|
||||
|
||||
// Unknown keys silently ignored (cluster, jetstream, gateway, leafnode, etc.)
|
||||
default:
|
||||
break;
|
||||
@@ -620,6 +626,145 @@ public static class ConfigProcessor
|
||||
opts.Tags = tags;
|
||||
}
|
||||
|
||||
// ─── MQTT parsing ────────────────────────────────────────────────
|
||||
// Reference: Go server/opts.go parseMQTT (lines ~5443-5541)
|
||||
|
||||
private static void ParseMqtt(Dictionary<string, object?> dict, NatsOptions opts, List<string> errors)
|
||||
{
|
||||
var mqtt = opts.Mqtt ?? new MqttOptions();
|
||||
|
||||
foreach (var (key, value) in dict)
|
||||
{
|
||||
switch (key.ToLowerInvariant())
|
||||
{
|
||||
case "listen":
|
||||
var (host, port) = ParseHostPort(value);
|
||||
if (host is not null) mqtt.Host = host;
|
||||
if (port is not null) mqtt.Port = port.Value;
|
||||
break;
|
||||
case "port":
|
||||
mqtt.Port = ToInt(value);
|
||||
break;
|
||||
case "host" or "net":
|
||||
mqtt.Host = ToString(value);
|
||||
break;
|
||||
case "no_auth_user":
|
||||
mqtt.NoAuthUser = ToString(value);
|
||||
break;
|
||||
case "tls":
|
||||
if (value is Dictionary<string, object?> tlsDict)
|
||||
ParseMqttTls(tlsDict, mqtt, errors);
|
||||
break;
|
||||
case "authorization" or "authentication":
|
||||
if (value is Dictionary<string, object?> authDict)
|
||||
ParseMqttAuth(authDict, mqtt, errors);
|
||||
break;
|
||||
case "ack_wait" or "ackwait":
|
||||
mqtt.AckWait = ParseDuration(value);
|
||||
break;
|
||||
case "js_api_timeout" or "api_timeout":
|
||||
mqtt.JsApiTimeout = ParseDuration(value);
|
||||
break;
|
||||
case "max_ack_pending" or "max_pending" or "max_inflight":
|
||||
var pending = ToInt(value);
|
||||
if (pending < 0 || pending > 0xFFFF)
|
||||
errors.Add($"mqtt max_ack_pending invalid value {pending}, should be in [0..{0xFFFF}] range");
|
||||
else
|
||||
mqtt.MaxAckPending = (ushort)pending;
|
||||
break;
|
||||
case "js_domain":
|
||||
mqtt.JsDomain = ToString(value);
|
||||
break;
|
||||
case "stream_replicas":
|
||||
mqtt.StreamReplicas = ToInt(value);
|
||||
break;
|
||||
case "consumer_replicas":
|
||||
mqtt.ConsumerReplicas = ToInt(value);
|
||||
break;
|
||||
case "consumer_memory_storage":
|
||||
mqtt.ConsumerMemoryStorage = ToBool(value);
|
||||
break;
|
||||
case "consumer_inactive_threshold" or "consumer_auto_cleanup":
|
||||
mqtt.ConsumerInactiveThreshold = ParseDuration(value);
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
opts.Mqtt = mqtt;
|
||||
}
|
||||
|
||||
private static void ParseMqttAuth(Dictionary<string, object?> dict, MqttOptions mqtt, List<string> errors)
|
||||
{
|
||||
foreach (var (key, value) in dict)
|
||||
{
|
||||
switch (key.ToLowerInvariant())
|
||||
{
|
||||
case "user" or "username":
|
||||
mqtt.Username = ToString(value);
|
||||
break;
|
||||
case "pass" or "password":
|
||||
mqtt.Password = ToString(value);
|
||||
break;
|
||||
case "token":
|
||||
mqtt.Token = ToString(value);
|
||||
break;
|
||||
case "timeout":
|
||||
mqtt.AuthTimeout = ToDouble(value);
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static void ParseMqttTls(Dictionary<string, object?> dict, MqttOptions mqtt, List<string> errors)
|
||||
{
|
||||
foreach (var (key, value) in dict)
|
||||
{
|
||||
switch (key.ToLowerInvariant())
|
||||
{
|
||||
case "cert_file":
|
||||
mqtt.TlsCert = ToString(value);
|
||||
break;
|
||||
case "key_file":
|
||||
mqtt.TlsKey = ToString(value);
|
||||
break;
|
||||
case "ca_file":
|
||||
mqtt.TlsCaCert = ToString(value);
|
||||
break;
|
||||
case "verify":
|
||||
mqtt.TlsVerify = ToBool(value);
|
||||
break;
|
||||
case "verify_and_map":
|
||||
var map = ToBool(value);
|
||||
mqtt.TlsMap = map;
|
||||
if (map) mqtt.TlsVerify = true;
|
||||
break;
|
||||
case "timeout":
|
||||
mqtt.TlsTimeout = ToDouble(value);
|
||||
break;
|
||||
case "pinned_certs":
|
||||
if (value is List<object?> pinnedList)
|
||||
{
|
||||
var certs = new HashSet<string>(StringComparer.OrdinalIgnoreCase);
|
||||
foreach (var item in pinnedList)
|
||||
{
|
||||
if (item is string s)
|
||||
certs.Add(s.ToLowerInvariant());
|
||||
}
|
||||
|
||||
mqtt.TlsPinnedCerts = certs;
|
||||
}
|
||||
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ─── Type conversion helpers ───────────────────────────────────
|
||||
|
||||
private static int ToInt(object? value) => value switch
|
||||
@@ -653,6 +798,15 @@ public static class ConfigProcessor
|
||||
_ => throw new FormatException($"Cannot convert {value?.GetType().Name ?? "null"} to string"),
|
||||
};
|
||||
|
||||
private static double ToDouble(object? value) => value switch
|
||||
{
|
||||
double d => d,
|
||||
long l => l,
|
||||
int i => i,
|
||||
string s when double.TryParse(s, NumberStyles.Float, CultureInfo.InvariantCulture, out var d) => d,
|
||||
_ => throw new FormatException($"Cannot convert {value?.GetType().Name ?? "null"} to double"),
|
||||
};
|
||||
|
||||
private static IReadOnlyList<string> ToStringList(object? value)
|
||||
{
|
||||
if (value is List<object?> list)
|
||||
|
||||
Reference in New Issue
Block a user