mirror of
https://github.com/l5yth/potato-mesh.git
synced 2026-08-10 10:52:55 +02:00
web: refactor 1/7 data processing (#772)
* web: refactor 1/7 data processing * web: close coverage gaps in data_processing submodules Bring every file under lib/potato_mesh/application/data_processing/ to 100% line coverage so codecov/patch passes on the 1/7 refactor PR. The gap was a relocation of pre-existing untested branches; closing them here keeps the subsequent refactor PRs in the series unblocked. * Add unit tests covering canonical sender/recipient overrides, reply_id/emoji updates on existing rows, and the rare INSERT ConstraintException recovery path inside +insert_message+. * Cover the non-canonical reporter and per-neighbour resolution branches in +insert_neighbors+. * Cover the SQLException rescue in +upsert_ingestor+, the fallback_num branch in +touch_node_last_seen+, the limit fallback in +read_json_body+, the unrecognised-type branch in +store_decrypted_payload+, the +power+ telemetry_type fallback, the default-coercion path in +resolve_numeric_metric+, and the numeric/bare-hex paths in +canonical_node_parts+ and +coerce_trace_node_id+. Drop dead code surfaced while pinning behaviour: * +clear_encrypted+ in +insert_message+ has been initialised to +false+ and never reassigned since #633 dropped the decrypted-text override; remove it and the four dependent branches. * The +rescue ArgumentError; nil+ tails in +identity.resolve_node_num+ and +traces.coerce_trace_node_id+ are unreachable because every +Integer(...)+ call inside is guarded by a regex pre-check. Add a comment to the +data_processing.rb+ shim explaining that the +require_relative+ list is ordered by dependency rather than alphabetically, addressing review nit #5.
This commit is contained in:
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,51 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Allowed values for the +telemetry_type+ discriminator column.
|
||||
VALID_TELEMETRY_TYPES = %w[device environment power air_quality].freeze
|
||||
|
||||
# Half-window (seconds) for the meshcore content-level message dedup
|
||||
# in +insert_message+ and the matching one-shot backfill. Set to
|
||||
# roughly 3× the observed relay-retransmit delta (~10 s) so genuine
|
||||
# clock skew across co-operating ingestors still collapses, while
|
||||
# rapid legitimate re-sends ("ack", "ok", "test") ≥30 s apart remain
|
||||
# distinct rows. See issue #756 and ``CONTRACTS.md`` for rationale.
|
||||
#
|
||||
# IMPORTANT: widening this value only takes effect at runtime — the
|
||||
# one-shot backfill in +PotatoMesh::App::Database+ is frozen at
|
||||
# +MESHCORE_CONTENT_DEDUP_BACKFILL_VERSION+. To re-sweep pre-existing
|
||||
# rows that newly fall within an expanded window, bump the backfill
|
||||
# version so the migration re-runs on the next deploy.
|
||||
MESHCORE_CONTENT_DEDUP_WINDOW_SECONDS = 30
|
||||
|
||||
# Coerce a Ruby boolean into a SQLite integer (1/0) while passing through
|
||||
# any other value unchanged. Used when writing boolean node fields.
|
||||
#
|
||||
# @param value [Boolean, Object] value to coerce.
|
||||
# @return [Integer, Object] 1, 0, or the original value.
|
||||
def coerce_bool(value)
|
||||
case value
|
||||
when true then 1
|
||||
when false then 0
|
||||
else value
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,273 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Decode and store decrypted payloads in domain-specific tables.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param message [Hash] original message payload.
|
||||
# @param packet_id [Integer] packet identifier for the message.
|
||||
# @param decrypted [Hash] decrypted payload metadata.
|
||||
# @param rx_time [Integer] receive time.
|
||||
# @param rx_iso [String] ISO 8601 receive timestamp.
|
||||
# @param from_id [String, nil] canonical sender identifier.
|
||||
# @param to_id [String, nil] destination identifier.
|
||||
# @param channel [Integer, nil] channel index.
|
||||
# @param portnum [Object, nil] port number identifier.
|
||||
# @param hop_limit [Integer, nil] hop limit value.
|
||||
# @param snr [Numeric, nil] signal-to-noise ratio.
|
||||
# @param rssi [Integer, nil] RSSI value.
|
||||
# @return [void]
|
||||
def store_decrypted_payload(
|
||||
db,
|
||||
message,
|
||||
packet_id,
|
||||
decrypted,
|
||||
rx_time:,
|
||||
rx_iso:,
|
||||
from_id:,
|
||||
to_id:,
|
||||
channel:,
|
||||
portnum:,
|
||||
hop_limit:,
|
||||
snr:,
|
||||
rssi:
|
||||
)
|
||||
payload_bytes = decrypted[:payload]
|
||||
return false unless payload_bytes
|
||||
|
||||
portnum_value = coerce_integer(portnum || decrypted[:portnum])
|
||||
return false unless portnum_value
|
||||
|
||||
payload_b64 = Base64.strict_encode64(payload_bytes)
|
||||
supported_ports = [3, 4, 67, 70, 71]
|
||||
return false unless supported_ports.include?(portnum_value)
|
||||
|
||||
decoded = PotatoMesh::App::Meshtastic::PayloadDecoder.decode(
|
||||
portnum: portnum_value,
|
||||
payload_b64: payload_b64,
|
||||
)
|
||||
return false unless decoded.is_a?(Hash)
|
||||
return false unless decoded["payload"].is_a?(Hash)
|
||||
|
||||
common_payload = {
|
||||
"id" => packet_id,
|
||||
"packet_id" => packet_id,
|
||||
"rx_time" => rx_time,
|
||||
"rx_iso" => rx_iso,
|
||||
"from_id" => from_id,
|
||||
"to_id" => to_id,
|
||||
"channel" => channel,
|
||||
"portnum" => portnum_value.to_s,
|
||||
"hop_limit" => hop_limit,
|
||||
"snr" => snr,
|
||||
"rssi" => rssi,
|
||||
"lora_freq" => coerce_integer(message["lora_freq"] || message["loraFrequency"]),
|
||||
"modem_preset" => string_or_nil(message["modem_preset"] || message["modemPreset"]),
|
||||
"payload_b64" => payload_b64,
|
||||
"ingestor" => string_or_nil(message["ingestor"]),
|
||||
}
|
||||
|
||||
case decoded["type"]
|
||||
when "POSITION_APP"
|
||||
payload = common_payload.merge("position" => decoded["payload"])
|
||||
insert_position(db, payload)
|
||||
debug_log(
|
||||
"Stored decrypted position payload",
|
||||
context: "data_processing.store_decrypted_payload",
|
||||
message_id: packet_id,
|
||||
portnum: portnum_value,
|
||||
)
|
||||
true
|
||||
when "NODEINFO_APP"
|
||||
node_payload = normalize_decrypted_nodeinfo_payload(decoded["payload"])
|
||||
return false unless valid_decrypted_nodeinfo_payload?(node_payload)
|
||||
|
||||
node_id = string_or_nil(node_payload["id"]) || from_id
|
||||
node_num = coerce_integer(node_payload["num"]) ||
|
||||
coerce_integer(message["from_num"]) ||
|
||||
resolve_node_num(from_id, message)
|
||||
node_id ||= format("!%08x", node_num & 0xFFFFFFFF) if node_num
|
||||
return false unless node_id
|
||||
|
||||
payload = node_payload.merge(
|
||||
"num" => node_num,
|
||||
"lastHeard" => coerce_integer(node_payload["lastHeard"] || node_payload["last_heard"]) || rx_time,
|
||||
"snr" => node_payload.key?("snr") ? node_payload["snr"] : snr,
|
||||
"lora_freq" => common_payload["lora_freq"],
|
||||
"modem_preset" => common_payload["modem_preset"],
|
||||
)
|
||||
upsert_node(db, node_id, payload)
|
||||
debug_log(
|
||||
"Stored decrypted node payload",
|
||||
context: "data_processing.store_decrypted_payload",
|
||||
message_id: packet_id,
|
||||
portnum: portnum_value,
|
||||
node_id: node_id,
|
||||
)
|
||||
true
|
||||
when "TELEMETRY_APP"
|
||||
payload = common_payload.merge("telemetry" => decoded["payload"])
|
||||
insert_telemetry(db, payload)
|
||||
debug_log(
|
||||
"Stored decrypted telemetry payload",
|
||||
context: "data_processing.store_decrypted_payload",
|
||||
message_id: packet_id,
|
||||
portnum: portnum_value,
|
||||
)
|
||||
true
|
||||
when "NEIGHBORINFO_APP"
|
||||
neighbor_payload = decoded["payload"]
|
||||
neighbors = neighbor_payload["neighbors"]
|
||||
neighbors = [] unless neighbors.is_a?(Array)
|
||||
normalized_neighbors = neighbors.map do |neighbor|
|
||||
next unless neighbor.is_a?(Hash)
|
||||
{
|
||||
"neighbor_id" => neighbor["node_id"] || neighbor["nodeId"] || neighbor["id"],
|
||||
"snr" => neighbor["snr"],
|
||||
"rx_time" => neighbor["last_rx_time"],
|
||||
}.compact
|
||||
end.compact
|
||||
return false if normalized_neighbors.empty?
|
||||
|
||||
payload = common_payload.merge(
|
||||
"node_id" => neighbor_payload["node_id"] || from_id,
|
||||
"neighbors" => normalized_neighbors,
|
||||
"node_broadcast_interval_secs" => neighbor_payload["node_broadcast_interval_secs"],
|
||||
"last_sent_by_id" => neighbor_payload["last_sent_by_id"],
|
||||
)
|
||||
insert_neighbors(db, payload)
|
||||
debug_log(
|
||||
"Stored decrypted neighbor payload",
|
||||
context: "data_processing.store_decrypted_payload",
|
||||
message_id: packet_id,
|
||||
portnum: portnum_value,
|
||||
)
|
||||
true
|
||||
when "TRACEROUTE_APP"
|
||||
route = decoded["payload"]["route"]
|
||||
route_back = decoded["payload"]["route_back"]
|
||||
hops = route.is_a?(Array) ? route : route_back.is_a?(Array) ? route_back : []
|
||||
dest = hops.last if hops.is_a?(Array) && !hops.empty?
|
||||
src_num = coerce_integer(message["from_num"]) || resolve_node_num(from_id, message)
|
||||
payload = common_payload.merge(
|
||||
"src" => src_num,
|
||||
"dest" => dest,
|
||||
"hops" => hops,
|
||||
)
|
||||
insert_trace(db, payload)
|
||||
debug_log(
|
||||
"Stored decrypted traceroute payload",
|
||||
context: "data_processing.store_decrypted_payload",
|
||||
message_id: packet_id,
|
||||
portnum: portnum_value,
|
||||
)
|
||||
true
|
||||
else
|
||||
false
|
||||
end
|
||||
end
|
||||
|
||||
# Validate decoded NodeInfo payloads before upserting node records.
|
||||
#
|
||||
# @param payload [Object] decoded payload candidate.
|
||||
# @return [Boolean] true when the payload resembles a Meshtastic NodeInfo.
|
||||
def valid_decrypted_nodeinfo_payload?(payload)
|
||||
return false unless payload.is_a?(Hash)
|
||||
return false if payload.empty?
|
||||
return false unless payload["user"].is_a?(Hash)
|
||||
|
||||
return false if payload.key?("position") && !payload["position"].is_a?(Hash)
|
||||
return false if payload.key?("deviceMetrics") && !payload["deviceMetrics"].is_a?(Hash)
|
||||
return false unless nodeinfo_user_has_identifying_fields?(payload["user"])
|
||||
|
||||
true
|
||||
end
|
||||
|
||||
# Normalize decoded NodeInfo payload keys for +upsert_node+ compatibility.
|
||||
#
|
||||
# The Python decoder preserves protobuf field names, so nested hashes may
|
||||
# use +snake_case+ keys that +upsert_node+ does not read.
|
||||
#
|
||||
# @param payload [Object] decoded NodeInfo payload.
|
||||
# @return [Hash] normalized payload hash.
|
||||
def normalize_decrypted_nodeinfo_payload(payload)
|
||||
return {} unless payload.is_a?(Hash)
|
||||
|
||||
user = payload["user"]
|
||||
normalized_user = user.is_a?(Hash) ? user.dup : nil
|
||||
if normalized_user
|
||||
normalized_user["shortName"] ||= normalized_user["short_name"]
|
||||
normalized_user["longName"] ||= normalized_user["long_name"]
|
||||
normalized_user["hwModel"] ||= normalized_user["hw_model"]
|
||||
normalized_user["publicKey"] ||= normalized_user["public_key"]
|
||||
normalized_user["isUnmessagable"] = normalized_user["is_unmessagable"] if normalized_user.key?("is_unmessagable")
|
||||
end
|
||||
|
||||
metrics = payload["deviceMetrics"] || payload["device_metrics"]
|
||||
normalized_metrics = metrics.is_a?(Hash) ? metrics.dup : nil
|
||||
if normalized_metrics
|
||||
normalized_metrics["batteryLevel"] ||= normalized_metrics["battery_level"]
|
||||
normalized_metrics["channelUtilization"] ||= normalized_metrics["channel_utilization"]
|
||||
normalized_metrics["airUtilTx"] ||= normalized_metrics["air_util_tx"]
|
||||
normalized_metrics["uptimeSeconds"] ||= normalized_metrics["uptime_seconds"]
|
||||
end
|
||||
|
||||
position = payload["position"]
|
||||
normalized_position = position.is_a?(Hash) ? position.dup : nil
|
||||
if normalized_position
|
||||
normalized_position["precisionBits"] ||= normalized_position["precision_bits"]
|
||||
normalized_position["locationSource"] ||= normalized_position["location_source"]
|
||||
end
|
||||
|
||||
normalized = payload.dup
|
||||
normalized["user"] = normalized_user if normalized_user
|
||||
normalized["deviceMetrics"] = normalized_metrics if normalized_metrics
|
||||
normalized["position"] = normalized_position if normalized_position
|
||||
normalized["lastHeard"] ||= normalized["last_heard"]
|
||||
normalized["hopsAway"] ||= normalized["hops_away"]
|
||||
normalized["isFavorite"] = normalized["is_favorite"] if normalized.key?("is_favorite")
|
||||
normalized["hwModel"] ||= normalized["hw_model"]
|
||||
normalized
|
||||
end
|
||||
|
||||
# Validate that a decoded NodeInfo user section contains identifying data.
|
||||
#
|
||||
# @param user [Hash] decoded NodeInfo user payload.
|
||||
# @return [Boolean] true when at least one identifying field is present.
|
||||
def nodeinfo_user_has_identifying_fields?(user)
|
||||
identifying_fields = [
|
||||
user["id"],
|
||||
user["shortName"],
|
||||
user["short_name"],
|
||||
user["longName"],
|
||||
user["long_name"],
|
||||
user["macaddr"],
|
||||
user["hwModel"],
|
||||
user["hw_model"],
|
||||
user["publicKey"],
|
||||
user["public_key"],
|
||||
]
|
||||
|
||||
identifying_fields.any? do |value|
|
||||
value.is_a?(String) ? !value.strip.empty? : !value.nil?
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,199 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Resolve the numeric representation of a node identifier from a packet payload.
|
||||
#
|
||||
# The +payload["num"]+ field may arrive as an Integer, a decimal string, or
|
||||
# a hexadecimal string (with or without an +0x+ prefix). When the field is
|
||||
# absent or ambiguous the method falls back to decoding the hex portion of
|
||||
# +node_id+.
|
||||
#
|
||||
# @param node_id [String, nil] canonical node identifier in +!xxxxxxxx+ form.
|
||||
# @param payload [Hash] inbound message payload that may carry a +num+ field.
|
||||
# @return [Integer, nil] resolved 32-bit node number or +nil+ when undecidable.
|
||||
def resolve_node_num(node_id, payload)
|
||||
raw = payload["num"]
|
||||
|
||||
case raw
|
||||
when Integer
|
||||
return raw
|
||||
when Numeric
|
||||
return raw.to_i
|
||||
when String
|
||||
trimmed = raw.strip
|
||||
return nil if trimmed.empty?
|
||||
return Integer(trimmed, 10) if trimmed.match?(/\A[0-9]+\z/)
|
||||
return Integer(trimmed.delete_prefix("0x").delete_prefix("0X"), 16) if trimmed.match?(/\A0[xX][0-9A-Fa-f]+\z/)
|
||||
if trimmed.match?(/\A[0-9A-Fa-f]+\z/)
|
||||
canonical = node_id.is_a?(String) ? node_id.strip : ""
|
||||
return Integer(trimmed, 16) if canonical.match?(/\A!?[0-9A-Fa-f]+\z/)
|
||||
end
|
||||
end
|
||||
|
||||
return nil unless node_id.is_a?(String)
|
||||
|
||||
hex = node_id.strip
|
||||
return nil if hex.empty?
|
||||
hex = hex.delete_prefix("!")
|
||||
return nil unless hex.match?(/\A[0-9A-Fa-f]+\z/)
|
||||
|
||||
Integer(hex, 16)
|
||||
end
|
||||
|
||||
# Derive the canonical triplet for a node reference.
|
||||
#
|
||||
# Accepts an Integer node number, a hex string with or without the +!+
|
||||
# sigil, a decimal numeric string, or a +0x+-prefixed hex string. A
|
||||
# +fallback_num+ may be provided when +node_ref+ is nil.
|
||||
#
|
||||
# @param node_ref [Integer, String, nil] raw node identifier from a packet.
|
||||
# @param fallback_num [Integer, nil] numeric fallback when +node_ref+ is nil.
|
||||
# @return [Array(String, Integer, String), nil] tuple of
|
||||
# +[canonical_id, node_num, short_id]+ or +nil+ when the reference cannot
|
||||
# be resolved. +canonical_id+ is prefixed with +!+ and zero-padded to
|
||||
# eight lowercase hex digits. +short_id+ is the upper-case last four
|
||||
# hex digits used for display.
|
||||
def canonical_node_parts(node_ref, fallback_num = nil)
|
||||
fallback = coerce_integer(fallback_num)
|
||||
|
||||
hex = nil
|
||||
num = nil
|
||||
|
||||
case node_ref
|
||||
when Integer
|
||||
num = node_ref
|
||||
when Numeric
|
||||
num = node_ref.to_i
|
||||
when String
|
||||
trimmed = node_ref.strip
|
||||
return nil if trimmed.empty?
|
||||
|
||||
if trimmed.start_with?("!")
|
||||
hex = trimmed.delete_prefix("!")
|
||||
elsif trimmed.match?(/\A0[xX][0-9A-Fa-f]+\z/)
|
||||
hex = trimmed[2..].to_s
|
||||
elsif trimmed.match?(/\A-?\d+\z/)
|
||||
num = trimmed.to_i
|
||||
elsif trimmed.match?(/\A[0-9A-Fa-f]+\z/)
|
||||
hex = trimmed
|
||||
else
|
||||
return nil
|
||||
end
|
||||
when nil
|
||||
num = fallback if fallback
|
||||
else
|
||||
return nil
|
||||
end
|
||||
|
||||
num ||= fallback if fallback
|
||||
|
||||
if hex
|
||||
begin
|
||||
num ||= Integer(hex, 16)
|
||||
rescue ArgumentError
|
||||
return nil
|
||||
end
|
||||
elsif num
|
||||
return nil if num.negative?
|
||||
hex = format("%08x", num & 0xFFFFFFFF)
|
||||
else
|
||||
return nil
|
||||
end
|
||||
|
||||
return nil if hex.nil? || hex.empty?
|
||||
|
||||
begin
|
||||
parsed = Integer(hex, 16)
|
||||
rescue ArgumentError
|
||||
return nil
|
||||
end
|
||||
|
||||
parsed &= 0xFFFFFFFF
|
||||
canonical_hex = format("%08x", parsed)
|
||||
short_id = canonical_hex[-4, 4].upcase
|
||||
|
||||
["!#{canonical_hex}", parsed, short_id]
|
||||
end
|
||||
|
||||
# Detect whether a node reference resolves to the broadcast address.
|
||||
#
|
||||
# @param node_ref [Integer, String, nil] raw node reference.
|
||||
# @param fallback_num [Integer, nil] optional numeric fallback.
|
||||
# @return [Boolean] true when the reference matches the broadcast address.
|
||||
def broadcast_node_ref?(node_ref, fallback_num = nil)
|
||||
return true if fallback_num == 0xFFFFFFFF
|
||||
trimmed = string_or_nil(node_ref)
|
||||
return false unless trimmed
|
||||
normalized = trimmed.delete_prefix("!").strip.downcase
|
||||
normalized == "ffffffff"
|
||||
end
|
||||
|
||||
# Converts a protocol identifier such as +meshtastic+ or +mesh-core+ into
|
||||
# the display label used in generated node names: capitalised parts joined
|
||||
# without a separator (e.g. +Meshtastic+, +MeshCore+).
|
||||
#
|
||||
# @param protocol [String] protocol identifier.
|
||||
# @return [String] formatted display label.
|
||||
def protocol_display_label(protocol)
|
||||
protocol.split(/[-_]/).map(&:capitalize).join
|
||||
end
|
||||
|
||||
# Returns true if +long_name+ is the synthetic placeholder generated by
|
||||
# +ensure_unknown_node+ for the given +node_id+ and +protocol+. Such
|
||||
# names carry no real information and must not overwrite a known name
|
||||
# already on record.
|
||||
#
|
||||
# @param long_name [String, nil] candidate long name.
|
||||
# @param node_id [String, nil] canonical node identifier.
|
||||
# @param protocol [String] protocol identifier the placeholder was generated for.
|
||||
# @return [Boolean] true when the long name is a generic placeholder.
|
||||
def generic_fallback_name?(long_name, node_id, protocol)
|
||||
return false unless long_name && !long_name.empty?
|
||||
|
||||
parts = canonical_node_parts(node_id)
|
||||
return false unless parts
|
||||
|
||||
short_id = parts[2]
|
||||
long_name == "#{protocol_display_label(protocol)} #{short_id}"
|
||||
end
|
||||
|
||||
# Resolve a raw node reference to its canonical row in the +nodes+ table.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param node_ref [Object] raw reference (string, integer, or hex string).
|
||||
# @return [String, nil] canonical +node_id+ or nil when no match exists.
|
||||
def normalize_node_id(db, node_ref)
|
||||
return nil if node_ref.nil?
|
||||
ref_str = node_ref.to_s.strip
|
||||
return nil if ref_str.empty?
|
||||
|
||||
node_id = db.get_first_value("SELECT node_id FROM nodes WHERE node_id = ?", [ref_str])
|
||||
return node_id if node_id
|
||||
|
||||
begin
|
||||
ref_num = Integer(ref_str, 10)
|
||||
rescue ArgumentError
|
||||
return nil
|
||||
end
|
||||
|
||||
db.get_first_value("SELECT node_id FROM nodes WHERE num = ?", [ref_num])
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,83 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Insert or update an ingestor heartbeat payload.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param payload [Hash] ingestor payload from the collector.
|
||||
# @return [Boolean] true when persistence succeeded.
|
||||
def upsert_ingestor(db, payload)
|
||||
return false unless payload.is_a?(Hash)
|
||||
|
||||
parts = canonical_node_parts(payload["node_id"] || payload["id"])
|
||||
return false unless parts
|
||||
|
||||
node_id, = parts
|
||||
now = Time.now.to_i
|
||||
|
||||
start_time = coerce_integer(payload["start_time"] || payload["startTime"]) || now
|
||||
last_seen_time =
|
||||
coerce_integer(payload["last_seen_time"] || payload["lastSeenTime"]) || start_time
|
||||
|
||||
start_time = 0 if start_time.negative?
|
||||
last_seen_time = 0 if last_seen_time.negative?
|
||||
start_time = now if start_time > now
|
||||
last_seen_time = now if last_seen_time > now
|
||||
last_seen_time = start_time if last_seen_time < start_time
|
||||
|
||||
version = string_or_nil(payload["version"] || payload["ingestorVersion"])
|
||||
return false unless version
|
||||
lora_freq = coerce_integer(payload["lora_freq"])
|
||||
modem_preset = string_or_nil(payload["modem_preset"])
|
||||
protocol = string_or_nil(payload["protocol"]) || "meshtastic"
|
||||
|
||||
with_busy_retry do
|
||||
db.execute <<~SQL, [node_id, start_time, last_seen_time, version, lora_freq, modem_preset, protocol]
|
||||
INSERT INTO ingestors(node_id, start_time, last_seen_time, version, lora_freq, modem_preset, protocol)
|
||||
VALUES(?,?,?,?,?,?,?)
|
||||
ON CONFLICT(node_id) DO UPDATE SET
|
||||
start_time = CASE
|
||||
WHEN excluded.start_time > ingestors.start_time THEN excluded.start_time
|
||||
ELSE ingestors.start_time
|
||||
END,
|
||||
last_seen_time = CASE
|
||||
WHEN excluded.last_seen_time > ingestors.last_seen_time THEN excluded.last_seen_time
|
||||
ELSE ingestors.last_seen_time
|
||||
END,
|
||||
version = COALESCE(excluded.version, ingestors.version),
|
||||
lora_freq = COALESCE(excluded.lora_freq, ingestors.lora_freq),
|
||||
modem_preset = COALESCE(excluded.modem_preset, ingestors.modem_preset),
|
||||
protocol = excluded.protocol
|
||||
SQL
|
||||
end
|
||||
|
||||
true
|
||||
rescue SQLite3::SQLException => e
|
||||
warn_log(
|
||||
"Failed to upsert ingestor record",
|
||||
context: "data_processing.ingestors",
|
||||
node_id: node_id,
|
||||
error_class: e.class.name,
|
||||
error_message: e.message,
|
||||
)
|
||||
false
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,494 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Determine whether the canonical sender identifier should override the
|
||||
# sender supplied by the ingestor. MeshCore packets that include a
|
||||
# +packet_id+ but no +id+ predate the canonical-id assignment, so we
|
||||
# prefer the canonical lookup when both are available.
|
||||
#
|
||||
# @param message [Hash] inbound message payload.
|
||||
# @return [Boolean] true when the canonical lookup wins.
|
||||
def prefer_canonical_sender?(message)
|
||||
message.is_a?(Hash) && message.key?("packet_id") && !message.key?("id")
|
||||
end
|
||||
|
||||
# Attempt to decrypt an encrypted Meshtastic message payload.
|
||||
#
|
||||
# @param message [Hash] message payload supplied by the ingestor.
|
||||
# @param packet_id [Integer] message packet identifier.
|
||||
# @param from_id [String, nil] canonical node identifier when available.
|
||||
# @param from_num [Integer, nil] numeric node identifier when available.
|
||||
# @param channel_index [Integer, nil] channel hash index.
|
||||
# @return [Hash, nil] decrypted payload metadata when parsing succeeds.
|
||||
def decrypt_meshtastic_message(message, packet_id, from_id, from_num, channel_index)
|
||||
return nil unless message.is_a?(Hash)
|
||||
|
||||
cipher_b64 = string_or_nil(message["encrypted"])
|
||||
return nil unless cipher_b64
|
||||
if (ENV["RACK_ENV"] == "test" || ENV["APP_ENV"] == "test" || defined?(RSpec)) &&
|
||||
ENV["MESHTASTIC_PSK_B64"].nil?
|
||||
return nil
|
||||
end
|
||||
|
||||
node_num = coerce_integer(from_num)
|
||||
if node_num.nil?
|
||||
parts = canonical_node_parts(from_id)
|
||||
node_num = parts[1] if parts
|
||||
end
|
||||
return nil unless node_num
|
||||
|
||||
psk_b64 = PotatoMesh::Config.meshtastic_psk_b64
|
||||
data = PotatoMesh::App::Meshtastic::Cipher.decrypt_data(
|
||||
cipher_b64: cipher_b64,
|
||||
packet_id: packet_id,
|
||||
from_id: from_id,
|
||||
from_num: node_num,
|
||||
psk_b64: psk_b64,
|
||||
)
|
||||
return nil unless data
|
||||
|
||||
channel_name = nil
|
||||
if channel_index.is_a?(Integer)
|
||||
candidates = PotatoMesh::App::Meshtastic::RainbowTable.channel_names_for(
|
||||
channel_index,
|
||||
psk_b64: psk_b64,
|
||||
)
|
||||
channel_name = candidates.first if candidates.any?
|
||||
end
|
||||
|
||||
{
|
||||
text: data[:text],
|
||||
portnum: data[:portnum],
|
||||
payload: data[:payload],
|
||||
channel_name: channel_name,
|
||||
}
|
||||
end
|
||||
|
||||
# Persist a chat-layer message payload, performing meshcore content
|
||||
# dedup, decryption, and per-protocol bookkeeping.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param message [Hash] inbound message payload.
|
||||
# @param protocol_cache [Hash, nil] optional per-batch ingestor protocol cache.
|
||||
# @return [void]
|
||||
def insert_message(db, message, protocol_cache: nil)
|
||||
return unless message.is_a?(Hash)
|
||||
|
||||
msg_id = coerce_integer(message["id"] || message["packet_id"])
|
||||
return unless msg_id
|
||||
|
||||
now = Time.now.to_i
|
||||
rx_time = coerce_integer(message["rx_time"])
|
||||
rx_time = now if rx_time.nil? || rx_time > now
|
||||
rx_iso = string_or_nil(message["rx_iso"])
|
||||
rx_iso ||= Time.at(rx_time).utc.iso8601
|
||||
|
||||
raw_from_id = message["from_id"]
|
||||
if raw_from_id.nil? || raw_from_id.to_s.strip.empty?
|
||||
alt_from = message["from"]
|
||||
raw_from_id = alt_from unless alt_from.nil? || alt_from.to_s.strip.empty?
|
||||
end
|
||||
|
||||
trimmed_from_id = string_or_nil(raw_from_id)
|
||||
canonical_from_id = string_or_nil(normalize_node_id(db, raw_from_id))
|
||||
from_id = trimmed_from_id
|
||||
if canonical_from_id
|
||||
if from_id.nil?
|
||||
from_id = canonical_from_id
|
||||
elsif prefer_canonical_sender?(message)
|
||||
from_id = canonical_from_id
|
||||
elsif from_id.start_with?("!") && from_id.casecmp(canonical_from_id) != 0
|
||||
from_id = canonical_from_id
|
||||
end
|
||||
end
|
||||
if from_id && !from_id.start_with?("^")
|
||||
canonical_parts = canonical_node_parts(from_id, message["from_num"])
|
||||
if canonical_parts && !from_id.start_with?("!")
|
||||
from_id = canonical_parts[0]
|
||||
message["from_num"] ||= canonical_parts[1]
|
||||
end
|
||||
end
|
||||
sender_present = !from_id.nil? || !coerce_integer(message["from_num"]).nil? || !trimmed_from_id.nil?
|
||||
|
||||
raw_to_id = message["to_id"]
|
||||
raw_to_id = message["to"] if raw_to_id.nil? || raw_to_id.to_s.strip.empty?
|
||||
trimmed_to_id = string_or_nil(raw_to_id)
|
||||
canonical_to_id = string_or_nil(normalize_node_id(db, raw_to_id))
|
||||
to_id = trimmed_to_id
|
||||
if canonical_to_id
|
||||
if to_id.nil?
|
||||
to_id = canonical_to_id
|
||||
elsif to_id.start_with?("!") && to_id.casecmp(canonical_to_id) != 0
|
||||
to_id = canonical_to_id
|
||||
end
|
||||
end
|
||||
if to_id && !to_id.start_with?("^")
|
||||
canonical_parts = canonical_node_parts(to_id, message["to_num"])
|
||||
if canonical_parts && !to_id.start_with?("!")
|
||||
to_id = canonical_parts[0]
|
||||
message["to_num"] ||= canonical_parts[1]
|
||||
end
|
||||
end
|
||||
|
||||
encrypted = string_or_nil(message["encrypted"])
|
||||
text = message["text"]
|
||||
portnum = message["portnum"]
|
||||
channel_index = coerce_integer(message["channel"] || message["channel_index"] || message["channelIndex"])
|
||||
|
||||
decrypted_payload = nil
|
||||
decrypted_portnum = nil
|
||||
|
||||
if encrypted && (text.nil? || text.to_s.strip.empty?)
|
||||
decrypted = decrypt_meshtastic_message(
|
||||
message,
|
||||
msg_id,
|
||||
from_id,
|
||||
message["from_num"],
|
||||
channel_index,
|
||||
)
|
||||
|
||||
if decrypted
|
||||
decrypted_payload = decrypted
|
||||
decrypted_portnum = decrypted[:portnum]
|
||||
end
|
||||
end
|
||||
|
||||
if encrypted && (text.nil? || text.to_s.strip.empty?)
|
||||
portnum = nil
|
||||
message.delete("portnum")
|
||||
end
|
||||
|
||||
lora_freq = coerce_integer(message["lora_freq"] || message["loraFrequency"])
|
||||
modem_preset = string_or_nil(message["modem_preset"] || message["modemPreset"])
|
||||
channel_name = string_or_nil(message["channel_name"] || message["channelName"])
|
||||
reply_id = coerce_integer(message["reply_id"] || message["replyId"])
|
||||
emoji = string_or_nil(message["emoji"])
|
||||
ingestor = string_or_nil(message["ingestor"])
|
||||
protocol = resolve_protocol(db, ingestor, cache: protocol_cache)
|
||||
|
||||
row = [
|
||||
msg_id,
|
||||
rx_time,
|
||||
rx_iso,
|
||||
from_id,
|
||||
to_id,
|
||||
message["channel"],
|
||||
portnum,
|
||||
text,
|
||||
encrypted,
|
||||
message["snr"],
|
||||
message["rssi"],
|
||||
message["hop_limit"],
|
||||
lora_freq,
|
||||
modem_preset,
|
||||
channel_name,
|
||||
reply_id,
|
||||
emoji,
|
||||
ingestor,
|
||||
protocol,
|
||||
]
|
||||
|
||||
with_busy_retry do
|
||||
# Meshcore-only content-level dedup (issue #756). The deterministic
|
||||
# message id (``_derive_message_id`` in the Python ingestor) hashes
|
||||
# ``sender_timestamp`` among other fields, but the MeshCore library
|
||||
# has been observed delivering the same physical packet twice with
|
||||
# a rewritten ``sender_timestamp`` (relay/retransmit behaviour).
|
||||
# The PK path below cannot catch that — two copies compute two
|
||||
# different ids — so we add a narrow content+window pre-check here.
|
||||
#
|
||||
# Ruby integer ``0`` is truthy, so the ``channel_index`` guard
|
||||
# passes for the broadcast channel intentionally; we only skip when
|
||||
# the channel is absent/nil. ``from_id`` + non-empty ``text`` keep
|
||||
# encrypted or anonymous traffic on the id-PK path.
|
||||
#
|
||||
# Known race: the SELECT and the downstream INSERT do not share a
|
||||
# transaction, so two Puma threads carrying the same content with
|
||||
# different ids can both pass the pre-check and both insert. The
|
||||
# deploy-time backfill sweeps the survivors; wrapping the pair in
|
||||
# ``db.transaction(:immediate)`` is a future tightening if the race
|
||||
# is ever observed in production.
|
||||
if protocol == "meshcore" && from_id && channel_index && text && !text.to_s.empty?
|
||||
# ``channel = ?`` matches the ``channel_index`` bind cleanly
|
||||
# because the guard above rejects nil; ``to_id`` may legitimately
|
||||
# be nil (rare meshcore fallback), so it keeps ``IS ?`` for a
|
||||
# NULL-safe compare.
|
||||
duplicate_id = db.get_first_value(
|
||||
<<~SQL,
|
||||
SELECT id FROM messages
|
||||
WHERE protocol = 'meshcore'
|
||||
AND from_id = ?
|
||||
AND to_id IS ?
|
||||
AND channel = ?
|
||||
AND text = ?
|
||||
AND rx_time BETWEEN ? AND ?
|
||||
AND id != ?
|
||||
LIMIT 1
|
||||
SQL
|
||||
[from_id, to_id, channel_index, text,
|
||||
rx_time - MESHCORE_CONTENT_DEDUP_WINDOW_SECONDS,
|
||||
rx_time + MESHCORE_CONTENT_DEDUP_WINDOW_SECONDS, msg_id],
|
||||
)
|
||||
if duplicate_id
|
||||
debug_log(
|
||||
"Skipped meshcore message duplicate",
|
||||
context: "data_processing.insert_message",
|
||||
new_id: msg_id,
|
||||
existing_id: duplicate_id,
|
||||
from_id: from_id,
|
||||
channel: channel_index,
|
||||
)
|
||||
return
|
||||
end
|
||||
end
|
||||
|
||||
existing = db.get_first_row(
|
||||
"SELECT from_id, to_id, text, encrypted, lora_freq, modem_preset, channel_name, reply_id, emoji, portnum, ingestor, protocol FROM messages WHERE id = ?",
|
||||
[msg_id],
|
||||
)
|
||||
if existing
|
||||
updates = {}
|
||||
existing_text = existing.is_a?(Hash) ? existing["text"] : existing[2]
|
||||
existing_text_str = existing_text&.to_s
|
||||
existing_has_text = existing_text_str && !existing_text_str.strip.empty?
|
||||
existing_from = existing.is_a?(Hash) ? existing["from_id"] : existing[0]
|
||||
existing_from_str = existing_from&.to_s
|
||||
return if !sender_present && (existing_from_str.nil? || existing_from_str.strip.empty?)
|
||||
existing_encrypted = existing.is_a?(Hash) ? existing["encrypted"] : existing[3]
|
||||
existing_encrypted_str = existing_encrypted&.to_s
|
||||
decrypted_precedence = text && existing_encrypted_str && !existing_encrypted_str.strip.empty?
|
||||
|
||||
if from_id
|
||||
should_update = existing_from_str.nil? || existing_from_str.strip.empty?
|
||||
should_update ||= existing_from != from_id
|
||||
updates["from_id"] = from_id if should_update
|
||||
end
|
||||
|
||||
if to_id
|
||||
existing_to = existing.is_a?(Hash) ? existing["to_id"] : existing[1]
|
||||
existing_to_str = existing_to&.to_s
|
||||
should_update = existing_to_str.nil? || existing_to_str.strip.empty?
|
||||
should_update ||= existing_to != to_id
|
||||
updates["to_id"] = to_id if should_update
|
||||
end
|
||||
|
||||
if decrypted_precedence && existing_encrypted_str && !existing_encrypted_str.strip.empty?
|
||||
updates["encrypted"] = nil if existing_encrypted
|
||||
elsif encrypted && !existing_has_text
|
||||
should_update = existing_encrypted_str.nil? || existing_encrypted_str.strip.empty?
|
||||
should_update ||= existing_encrypted != encrypted
|
||||
updates["encrypted"] = encrypted if should_update
|
||||
end
|
||||
|
||||
if text
|
||||
should_update = existing_text_str.nil? || existing_text_str.strip.empty?
|
||||
should_update ||= existing_text != text
|
||||
updates["text"] = text if should_update
|
||||
end
|
||||
|
||||
if decrypted_precedence
|
||||
updates["channel"] = message["channel"] if message.key?("channel")
|
||||
updates["snr"] = message["snr"] if message.key?("snr")
|
||||
updates["rssi"] = message["rssi"] if message.key?("rssi")
|
||||
updates["hop_limit"] = message["hop_limit"] if message.key?("hop_limit")
|
||||
updates["lora_freq"] = lora_freq unless lora_freq.nil?
|
||||
updates["modem_preset"] = modem_preset if modem_preset
|
||||
updates["channel_name"] = channel_name if channel_name
|
||||
updates["rx_time"] = rx_time if rx_time
|
||||
updates["rx_iso"] = rx_iso if rx_iso
|
||||
end
|
||||
|
||||
if portnum
|
||||
existing_portnum = existing.is_a?(Hash) ? existing["portnum"] : existing[9]
|
||||
existing_portnum_str = existing_portnum&.to_s
|
||||
should_update = existing_portnum_str.nil? || existing_portnum_str.strip.empty?
|
||||
should_update ||= existing_portnum != portnum
|
||||
should_update ||= decrypted_precedence
|
||||
updates["portnum"] = portnum if should_update
|
||||
end
|
||||
|
||||
unless lora_freq.nil?
|
||||
existing_lora = existing.is_a?(Hash) ? existing["lora_freq"] : existing[4]
|
||||
updates["lora_freq"] = lora_freq if existing_lora != lora_freq
|
||||
end
|
||||
|
||||
if modem_preset
|
||||
existing_preset = existing.is_a?(Hash) ? existing["modem_preset"] : existing[5]
|
||||
existing_preset_str = existing_preset&.to_s
|
||||
should_update = existing_preset_str.nil? || existing_preset_str.strip.empty?
|
||||
should_update ||= existing_preset != modem_preset
|
||||
updates["modem_preset"] = modem_preset if should_update
|
||||
end
|
||||
|
||||
if channel_name
|
||||
existing_channel = existing.is_a?(Hash) ? existing["channel_name"] : existing[6]
|
||||
existing_channel_str = existing_channel&.to_s
|
||||
should_update = existing_channel_str.nil? || existing_channel_str.strip.empty?
|
||||
should_update ||= existing_channel != channel_name
|
||||
updates["channel_name"] = channel_name if should_update
|
||||
end
|
||||
|
||||
unless reply_id.nil?
|
||||
existing_reply = existing.is_a?(Hash) ? existing["reply_id"] : existing[7]
|
||||
updates["reply_id"] = reply_id if existing_reply != reply_id
|
||||
end
|
||||
|
||||
if emoji
|
||||
existing_emoji = existing.is_a?(Hash) ? existing["emoji"] : existing[8]
|
||||
existing_emoji_str = existing_emoji&.to_s
|
||||
should_update = existing_emoji_str.nil? || existing_emoji_str.strip.empty?
|
||||
should_update ||= existing_emoji != emoji
|
||||
updates["emoji"] = emoji if should_update
|
||||
end
|
||||
|
||||
if ingestor
|
||||
existing_ingestor = existing.is_a?(Hash) ? existing["ingestor"] : existing[10]
|
||||
existing_ingestor = string_or_nil(existing_ingestor)
|
||||
updates["ingestor"] = ingestor if existing_ingestor.nil?
|
||||
end
|
||||
|
||||
existing_protocol = existing.is_a?(Hash) ? existing["protocol"] : existing[11]
|
||||
return if existing_protocol && existing_protocol != "meshtastic" && existing_protocol != protocol
|
||||
updates["protocol"] = protocol if (existing_protocol.nil? || existing_protocol == "meshtastic") && protocol != "meshtastic"
|
||||
|
||||
unless updates.empty?
|
||||
assignments = updates.keys.map { |column| "#{column} = ?" }.join(", ")
|
||||
db.execute("UPDATE messages SET #{assignments} WHERE id = ?", updates.values + [msg_id])
|
||||
end
|
||||
else
|
||||
PotatoMesh::App::Prometheus::MESSAGES_TOTAL.increment
|
||||
|
||||
begin
|
||||
db.execute <<~SQL, row
|
||||
INSERT INTO messages(id,rx_time,rx_iso,from_id,to_id,channel,portnum,text,encrypted,snr,rssi,hop_limit,lora_freq,modem_preset,channel_name,reply_id,emoji,ingestor,protocol)
|
||||
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
|
||||
SQL
|
||||
rescue SQLite3::ConstraintException
|
||||
existing_row = db.get_first_row(
|
||||
"SELECT text, encrypted, ingestor, protocol FROM messages WHERE id = ?",
|
||||
[msg_id],
|
||||
)
|
||||
existing_text = existing_row.is_a?(Hash) ? existing_row["text"] : existing_row&.[](0)
|
||||
existing_text_str = existing_text&.to_s
|
||||
allow_encrypted_update = existing_text_str.nil? || existing_text_str.strip.empty?
|
||||
existing_encrypted = existing_row.is_a?(Hash) ? existing_row["encrypted"] : existing_row&.[](1)
|
||||
existing_encrypted_str = existing_encrypted&.to_s
|
||||
existing_ingestor = existing_row.is_a?(Hash) ? existing_row["ingestor"] : existing_row&.[](2)
|
||||
existing_ingestor = string_or_nil(existing_ingestor)
|
||||
existing_fallback_protocol = existing_row.is_a?(Hash) ? existing_row["protocol"] : existing_row&.[](3)
|
||||
# Guard against cross-protocol contamination in the constraint fallback path,
|
||||
# mirroring the same guard applied in the primary update path above.
|
||||
return if existing_fallback_protocol && existing_fallback_protocol != "meshtastic" && existing_fallback_protocol != protocol
|
||||
decrypted_precedence = text && existing_encrypted_str && !existing_encrypted_str.strip.empty?
|
||||
|
||||
fallback_updates = {}
|
||||
fallback_updates["from_id"] = from_id if from_id
|
||||
fallback_updates["to_id"] = to_id if to_id
|
||||
fallback_updates["text"] = text if text
|
||||
fallback_updates["encrypted"] = encrypted if encrypted && allow_encrypted_update
|
||||
fallback_updates["portnum"] = portnum if portnum
|
||||
if decrypted_precedence
|
||||
fallback_updates["channel"] = message["channel"] if message.key?("channel")
|
||||
fallback_updates["snr"] = message["snr"] if message.key?("snr")
|
||||
fallback_updates["rssi"] = message["rssi"] if message.key?("rssi")
|
||||
fallback_updates["hop_limit"] = message["hop_limit"] if message.key?("hop_limit")
|
||||
fallback_updates["portnum"] = portnum if portnum
|
||||
fallback_updates["lora_freq"] = lora_freq unless lora_freq.nil?
|
||||
fallback_updates["modem_preset"] = modem_preset if modem_preset
|
||||
fallback_updates["channel_name"] = channel_name if channel_name
|
||||
fallback_updates["rx_time"] = rx_time if rx_time
|
||||
fallback_updates["rx_iso"] = rx_iso if rx_iso
|
||||
else
|
||||
fallback_updates["lora_freq"] = lora_freq unless lora_freq.nil?
|
||||
fallback_updates["modem_preset"] = modem_preset if modem_preset
|
||||
fallback_updates["channel_name"] = channel_name if channel_name
|
||||
end
|
||||
fallback_updates["reply_id"] = reply_id unless reply_id.nil?
|
||||
fallback_updates["emoji"] = emoji if emoji
|
||||
fallback_updates["ingestor"] = ingestor if ingestor && existing_ingestor.nil?
|
||||
fallback_updates["protocol"] = protocol if (existing_fallback_protocol.nil? || existing_fallback_protocol == "meshtastic") && protocol != "meshtastic"
|
||||
unless fallback_updates.empty?
|
||||
assignments = fallback_updates.keys.map { |column| "#{column} = ?" }.join(", ")
|
||||
db.execute("UPDATE messages SET #{assignments} WHERE id = ?", fallback_updates.values + [msg_id])
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
stored_decrypted = nil
|
||||
if decrypted_payload
|
||||
stored_decrypted = store_decrypted_payload(
|
||||
db,
|
||||
message,
|
||||
msg_id,
|
||||
decrypted_payload,
|
||||
rx_time: rx_time,
|
||||
rx_iso: rx_iso,
|
||||
from_id: from_id,
|
||||
to_id: to_id,
|
||||
channel: message["channel"],
|
||||
portnum: portnum || decrypted_portnum,
|
||||
hop_limit: message["hop_limit"],
|
||||
snr: message["snr"],
|
||||
rssi: message["rssi"],
|
||||
)
|
||||
end
|
||||
|
||||
if stored_decrypted && encrypted
|
||||
with_busy_retry do
|
||||
db.execute("UPDATE messages SET encrypted = NULL WHERE id = ?", [msg_id])
|
||||
end
|
||||
debug_log(
|
||||
"Cleared encrypted payload after decoding",
|
||||
context: "data_processing.insert_message",
|
||||
message_id: msg_id,
|
||||
portnum: portnum || decrypted_portnum,
|
||||
)
|
||||
end
|
||||
|
||||
should_touch_message = !stored_decrypted
|
||||
if should_touch_message
|
||||
ensure_unknown_node(db, from_id || raw_from_id, message["from_num"], heard_time: rx_time, protocol: protocol)
|
||||
touch_node_last_seen(
|
||||
db,
|
||||
from_id || raw_from_id || message["from_num"],
|
||||
message["from_num"],
|
||||
rx_time: rx_time,
|
||||
source: :message,
|
||||
lora_freq: lora_freq,
|
||||
modem_preset: modem_preset,
|
||||
)
|
||||
|
||||
ensure_unknown_node(db, to_id || raw_to_id, message["to_num"], heard_time: rx_time, protocol: protocol) if to_id || raw_to_id
|
||||
if to_id || raw_to_id || message.key?("to_num")
|
||||
touch_node_last_seen(
|
||||
db,
|
||||
to_id || raw_to_id || message["to_num"],
|
||||
message["to_num"],
|
||||
rx_time: rx_time,
|
||||
source: :message,
|
||||
lora_freq: lora_freq,
|
||||
modem_preset: modem_preset,
|
||||
)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,144 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Persist a neighbours snapshot for a single reporting node.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param payload [Hash] inbound NeighborInfo payload.
|
||||
# @param protocol_cache [Hash, nil] optional per-batch ingestor protocol cache.
|
||||
# @return [void]
|
||||
def insert_neighbors(db, payload, protocol_cache: nil)
|
||||
return unless payload.is_a?(Hash)
|
||||
|
||||
now = Time.now.to_i
|
||||
rx_time = coerce_integer(payload["rx_time"])
|
||||
rx_time = now if rx_time.nil? || rx_time > now
|
||||
|
||||
raw_node_id = payload["node_id"] || payload["node"] || payload["from_id"]
|
||||
raw_node_num = coerce_integer(payload["node_num"]) || coerce_integer(payload["num"])
|
||||
|
||||
canonical_parts = canonical_node_parts(raw_node_id, raw_node_num)
|
||||
if canonical_parts
|
||||
node_id, node_num, = canonical_parts
|
||||
else
|
||||
node_id = string_or_nil(raw_node_id)
|
||||
canonical = normalize_node_id(db, node_id || raw_node_num)
|
||||
node_id = canonical if canonical
|
||||
if node_id&.start_with?("!") && raw_node_num.nil?
|
||||
begin
|
||||
node_num = Integer(node_id.delete_prefix("!"), 16)
|
||||
rescue ArgumentError
|
||||
node_num = nil
|
||||
end
|
||||
else
|
||||
node_num = raw_node_num
|
||||
end
|
||||
end
|
||||
|
||||
return unless node_id
|
||||
|
||||
node_id = "!#{node_id.delete_prefix("!").downcase}" if node_id.start_with?("!")
|
||||
|
||||
ingestor = string_or_nil(payload["ingestor"])
|
||||
protocol = resolve_protocol(db, ingestor, cache: protocol_cache)
|
||||
|
||||
ensure_unknown_node(db, node_id || node_num, node_num, heard_time: rx_time, protocol: protocol)
|
||||
touch_node_last_seen(db, node_id || node_num, node_num, rx_time: rx_time, source: :neighborinfo)
|
||||
|
||||
neighbor_entries = []
|
||||
neighbors_payload = payload["neighbors"]
|
||||
neighbors_list = neighbors_payload.is_a?(Array) ? neighbors_payload : []
|
||||
|
||||
neighbors_list.each do |neighbor|
|
||||
next unless neighbor.is_a?(Hash)
|
||||
|
||||
neighbor_ref = neighbor["neighbor_id"] || neighbor["node_id"] || neighbor["nodeId"] || neighbor["id"]
|
||||
neighbor_num = coerce_integer(
|
||||
neighbor["neighbor_num"] || neighbor["node_num"] || neighbor["nodeId"] || neighbor["id"],
|
||||
)
|
||||
|
||||
canonical_neighbor = canonical_node_parts(neighbor_ref, neighbor_num)
|
||||
if canonical_neighbor
|
||||
neighbor_id, neighbor_num, = canonical_neighbor
|
||||
else
|
||||
neighbor_id = string_or_nil(neighbor_ref)
|
||||
canonical_neighbor_id = normalize_node_id(db, neighbor_id || neighbor_num)
|
||||
neighbor_id = canonical_neighbor_id if canonical_neighbor_id
|
||||
if neighbor_id&.start_with?("!") && neighbor_num.nil?
|
||||
begin
|
||||
neighbor_num = Integer(neighbor_id.delete_prefix("!"), 16)
|
||||
rescue ArgumentError
|
||||
neighbor_num = nil
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
next unless neighbor_id
|
||||
|
||||
neighbor_id = "!#{neighbor_id.delete_prefix("!").downcase}" if neighbor_id.start_with?("!")
|
||||
|
||||
entry_rx_time = coerce_integer(neighbor["rx_time"]) || rx_time
|
||||
entry_rx_time = now if entry_rx_time && entry_rx_time > now
|
||||
snr = coerce_float(neighbor["snr"])
|
||||
|
||||
ensure_unknown_node(db, neighbor_id || neighbor_num, neighbor_num, heard_time: entry_rx_time, protocol: protocol)
|
||||
|
||||
neighbor_entries << [neighbor_id, snr, entry_rx_time, ingestor, protocol]
|
||||
end
|
||||
|
||||
with_busy_retry do
|
||||
db.transaction do
|
||||
if neighbor_entries.empty?
|
||||
db.execute("DELETE FROM neighbors WHERE node_id = ?", [node_id])
|
||||
else
|
||||
expected_neighbors = neighbor_entries.map(&:first).uniq
|
||||
existing_neighbors = db.execute(
|
||||
"SELECT neighbor_id FROM neighbors WHERE node_id = ?",
|
||||
[node_id],
|
||||
).flatten
|
||||
stale_neighbors = existing_neighbors - expected_neighbors
|
||||
stale_neighbors.each_slice(500) do |slice|
|
||||
placeholders = slice.map { "?" }.join(",")
|
||||
db.execute(
|
||||
"DELETE FROM neighbors WHERE node_id = ? AND neighbor_id IN (#{placeholders})",
|
||||
[node_id] + slice,
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
neighbor_entries.each do |neighbor_id, snr_value, heard_time, reporter_id, proto|
|
||||
db.execute(
|
||||
<<~SQL,
|
||||
INSERT INTO neighbors(node_id, neighbor_id, snr, rx_time, ingestor, protocol)
|
||||
VALUES (?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(node_id, neighbor_id) DO UPDATE SET
|
||||
snr = excluded.snr,
|
||||
rx_time = excluded.rx_time,
|
||||
ingestor = COALESCE(NULLIF(neighbors.ingestor,''), excluded.ingestor),
|
||||
protocol = COALESCE(NULLIF(neighbors.protocol,'meshtastic'), excluded.protocol)
|
||||
SQL
|
||||
[node_id, neighbor_id, snr_value, heard_time, reporter_id, proto],
|
||||
)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,570 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Insert a hidden placeholder node when an unknown reference is encountered.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param node_ref [Object] raw node reference from the inbound payload.
|
||||
# @param fallback_num [Integer, nil] numeric fallback when +node_ref+ is nil.
|
||||
# @param heard_time [Integer, nil] timestamp to record as +last_heard+/+first_heard+.
|
||||
# @param protocol [String] protocol identifier for placeholder generation.
|
||||
# @return [Boolean, nil] true when a row was inserted, false/nil otherwise.
|
||||
def ensure_unknown_node(db, node_ref, fallback_num = nil, heard_time: nil, protocol: "meshtastic")
|
||||
parts = canonical_node_parts(node_ref, fallback_num)
|
||||
return unless parts
|
||||
|
||||
node_id, node_num, short_id = parts
|
||||
return if broadcast_node_ref?(node_id, node_num)
|
||||
|
||||
existing = db.get_first_value(
|
||||
"SELECT 1 FROM nodes WHERE node_id = ? LIMIT 1",
|
||||
[node_id],
|
||||
)
|
||||
return if existing
|
||||
|
||||
long_name = "#{protocol_display_label(protocol)} #{short_id}"
|
||||
default_role = case protocol
|
||||
when "meshcore" then "COMPANION"
|
||||
else "CLIENT_HIDDEN"
|
||||
end
|
||||
heard_time = coerce_integer(heard_time)
|
||||
inserted = false
|
||||
|
||||
with_busy_retry do
|
||||
db.execute(
|
||||
<<~SQL,
|
||||
INSERT OR IGNORE INTO nodes(node_id,num,short_name,long_name,role,last_heard,first_heard,protocol)
|
||||
VALUES (?,?,?,?,?,?,?,?)
|
||||
SQL
|
||||
[node_id, node_num, short_id, long_name, default_role, heard_time, heard_time, protocol],
|
||||
)
|
||||
inserted = db.changes.positive?
|
||||
end
|
||||
|
||||
if inserted
|
||||
debug_log(
|
||||
"Created hidden placeholder node",
|
||||
context: "data_processing.ensure_unknown_node",
|
||||
node_id: node_id,
|
||||
reference: node_ref,
|
||||
fallback: fallback_num,
|
||||
heard_time: heard_time,
|
||||
)
|
||||
end
|
||||
|
||||
inserted
|
||||
end
|
||||
|
||||
# Refresh a node's +last_heard+, +first_heard+, +lora_freq+, and
|
||||
# +modem_preset+ columns from a freshly received packet.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param node_ref [Object] raw node reference.
|
||||
# @param fallback_num [Integer, nil] numeric fallback when +node_ref+ is nil.
|
||||
# @param rx_time [Integer, nil] receive timestamp; the method exits early when nil.
|
||||
# @param source [Symbol, nil] originating subsystem (used for debug logs).
|
||||
# @param lora_freq [Integer, nil] LoRa frequency; only updated when non-nil.
|
||||
# @param modem_preset [String, nil] modem preset name; only updated when non-nil.
|
||||
# @return [Boolean] true when at least one row was updated.
|
||||
def touch_node_last_seen(
|
||||
db,
|
||||
node_ref,
|
||||
fallback_num = nil,
|
||||
rx_time: nil,
|
||||
source: nil,
|
||||
lora_freq: nil,
|
||||
modem_preset: nil
|
||||
)
|
||||
timestamp = coerce_integer(rx_time)
|
||||
return unless timestamp
|
||||
|
||||
node_id = nil
|
||||
|
||||
parts = canonical_node_parts(node_ref, fallback_num)
|
||||
if parts
|
||||
node_id, node_num = parts
|
||||
return if broadcast_node_ref?(node_id, node_num)
|
||||
end
|
||||
|
||||
unless node_id
|
||||
trimmed = string_or_nil(node_ref)
|
||||
if trimmed
|
||||
node_id = normalize_node_id(db, trimmed) || trimmed
|
||||
elsif fallback_num
|
||||
fallback_parts = canonical_node_parts(fallback_num, nil)
|
||||
node_id, = fallback_parts if fallback_parts
|
||||
end
|
||||
end
|
||||
|
||||
return if broadcast_node_ref?(node_id, fallback_num)
|
||||
return unless node_id
|
||||
|
||||
lora_freq = coerce_integer(lora_freq)
|
||||
modem_preset = string_or_nil(modem_preset)
|
||||
updated = false
|
||||
with_busy_retry do
|
||||
db.execute <<~SQL, [timestamp, timestamp, timestamp, lora_freq, modem_preset, node_id]
|
||||
UPDATE nodes
|
||||
SET last_heard = CASE
|
||||
WHEN COALESCE(last_heard, 0) >= ? THEN last_heard
|
||||
ELSE ?
|
||||
END,
|
||||
first_heard = COALESCE(first_heard, ?),
|
||||
lora_freq = COALESCE(?, lora_freq),
|
||||
modem_preset = COALESCE(?, modem_preset)
|
||||
WHERE node_id = ?
|
||||
SQL
|
||||
updated ||= db.changes.positive?
|
||||
end
|
||||
|
||||
if updated
|
||||
debug_log(
|
||||
"Updated node last seen timestamp",
|
||||
context: "data_processing.touch_node_last_seen",
|
||||
node_id: node_id,
|
||||
timestamp: timestamp,
|
||||
source: source || :unknown,
|
||||
lora_freq: lora_freq,
|
||||
modem_preset: modem_preset,
|
||||
)
|
||||
end
|
||||
|
||||
updated
|
||||
end
|
||||
|
||||
# Insert or update a node row from an inbound NodeInfo-style payload.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param node_id [String] canonical node identifier.
|
||||
# @param n [Hash] node payload extracted from the ingestor.
|
||||
# @param protocol [String] protocol identifier (default +meshtastic+).
|
||||
# @return [void]
|
||||
def upsert_node(db, node_id, n, protocol: "meshtastic")
|
||||
user = n["user"] || {}
|
||||
met = n["deviceMetrics"] || {}
|
||||
pos = n["position"] || {}
|
||||
# nil when user info absent; COALESCE in the conflict clause preserves
|
||||
# the stored role rather than overwriting with a default.
|
||||
role = user["role"]
|
||||
lh = coerce_integer(n["lastHeard"])
|
||||
pt = coerce_integer(pos["time"])
|
||||
now = Time.now.to_i
|
||||
pt = nil if pt && pt > now
|
||||
lh = now if lh && lh > now
|
||||
# 0 is truthy in Ruby — `lh ||= now` won't replace it, leaving the
|
||||
# 7-day list filter to evaluate `0 >= now-7days` → false (node hidden).
|
||||
lh = nil if lh && lh <= 0
|
||||
# position.time = 0 means no GPS fix; skip it as a last_heard anchor
|
||||
# (would re-introduce the same zero-timestamp exclusion bug for lh).
|
||||
lh = pt if pt && pt > 0 && (!lh || lh < pt)
|
||||
lh ||= now
|
||||
node_num = resolve_node_num(node_id, n)
|
||||
|
||||
update_prometheus_metrics(node_id, user, role, met, pos)
|
||||
|
||||
lora_freq = coerce_integer(n["lora_freq"] || n["loraFrequency"])
|
||||
modem_preset = string_or_nil(n["modem_preset"] || n["modemPreset"])
|
||||
# Synthetic flag: true for placeholder nodes created from channel message
|
||||
# sender names before the real contact advertisement is received.
|
||||
synthetic = user["synthetic"] ? 1 : 0
|
||||
long_name = user["longName"]
|
||||
|
||||
# If the incoming long name is a generic placeholder, prefer any real
|
||||
# name already on record so we never stomp known data with fallback
|
||||
# text. For new nodes there is nothing to preserve, so the generic
|
||||
# name is still written via the INSERT VALUES path.
|
||||
long_name_conflict_sql = if generic_fallback_name?(long_name, node_id, protocol)
|
||||
# Generic placeholder: keep any real name already on record.
|
||||
# COALESCE returns nodes.long_name when non-null, otherwise falls
|
||||
# back to the incoming generic — so brand-new nodes still get it.
|
||||
"COALESCE(nodes.long_name, excluded.long_name)"
|
||||
else
|
||||
# Real name (or nil): use the incoming value, preserving the
|
||||
# existing name only when the incoming value is nil. A nil
|
||||
# long_name in the packet carries no information, so falling back
|
||||
# to what we already have is better than overwriting with NULL.
|
||||
"COALESCE(excluded.long_name, nodes.long_name)"
|
||||
end
|
||||
|
||||
row = [
|
||||
node_id,
|
||||
node_num,
|
||||
user["shortName"],
|
||||
long_name,
|
||||
user["macaddr"],
|
||||
user["hwModel"] || n["hwModel"],
|
||||
role,
|
||||
user["publicKey"],
|
||||
coerce_bool(user["isUnmessagable"]),
|
||||
coerce_bool(n["isFavorite"]),
|
||||
n["hopsAway"],
|
||||
n["snr"],
|
||||
lh,
|
||||
lh,
|
||||
met["batteryLevel"],
|
||||
met["voltage"],
|
||||
met["channelUtilization"],
|
||||
met["airUtilTx"],
|
||||
met["uptimeSeconds"],
|
||||
pt,
|
||||
pos["locationSource"],
|
||||
coerce_integer(
|
||||
pos["precisionBits"] ||
|
||||
pos["precision_bits"] ||
|
||||
pos.dig("raw", "precision_bits"),
|
||||
),
|
||||
pos["latitude"],
|
||||
pos["longitude"],
|
||||
pos["altitude"],
|
||||
lora_freq,
|
||||
modem_preset,
|
||||
protocol,
|
||||
synthetic,
|
||||
]
|
||||
with_busy_retry do
|
||||
db.transaction do
|
||||
db.execute(<<~SQL, row)
|
||||
INSERT INTO nodes(node_id,num,short_name,long_name,macaddr,hw_model,role,public_key,is_unmessagable,is_favorite,
|
||||
hops_away,snr,last_heard,first_heard,battery_level,voltage,channel_utilization,air_util_tx,uptime_seconds,
|
||||
position_time,location_source,precision_bits,latitude,longitude,altitude,lora_freq,modem_preset,protocol,synthetic)
|
||||
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
|
||||
ON CONFLICT(node_id) DO UPDATE SET
|
||||
num=COALESCE(excluded.num, nodes.num),
|
||||
short_name=COALESCE(excluded.short_name, nodes.short_name),
|
||||
long_name=#{long_name_conflict_sql},
|
||||
macaddr=COALESCE(excluded.macaddr, nodes.macaddr),
|
||||
hw_model=COALESCE(excluded.hw_model, nodes.hw_model),
|
||||
role=COALESCE(excluded.role, nodes.role),
|
||||
public_key=COALESCE(excluded.public_key, nodes.public_key),
|
||||
is_unmessagable=COALESCE(excluded.is_unmessagable, nodes.is_unmessagable),
|
||||
is_favorite=excluded.is_favorite, hops_away=excluded.hops_away, snr=excluded.snr, last_heard=excluded.last_heard,
|
||||
first_heard=COALESCE(nodes.first_heard, excluded.first_heard, excluded.last_heard),
|
||||
battery_level=excluded.battery_level, voltage=excluded.voltage, channel_utilization=excluded.channel_utilization,
|
||||
air_util_tx=excluded.air_util_tx, uptime_seconds=excluded.uptime_seconds,
|
||||
position_time=COALESCE(excluded.position_time, nodes.position_time),
|
||||
location_source=COALESCE(excluded.location_source, nodes.location_source),
|
||||
precision_bits=COALESCE(excluded.precision_bits, nodes.precision_bits),
|
||||
latitude=COALESCE(excluded.latitude, nodes.latitude),
|
||||
longitude=COALESCE(excluded.longitude, nodes.longitude),
|
||||
altitude=COALESCE(excluded.altitude, nodes.altitude),
|
||||
lora_freq=excluded.lora_freq, modem_preset=excluded.modem_preset,
|
||||
protocol=COALESCE(NULLIF(nodes.protocol,'meshtastic'), excluded.protocol),
|
||||
synthetic=MIN(COALESCE(excluded.synthetic,1), COALESCE(nodes.synthetic,1))
|
||||
WHERE COALESCE(excluded.last_heard,0) >= COALESCE(nodes.last_heard,0)
|
||||
AND NOT (COALESCE(nodes.synthetic,0) = 0 AND excluded.synthetic = 1)
|
||||
SQL
|
||||
|
||||
# Reconcile synthetic placeholder rows with their real counterparts
|
||||
# whenever a MeshCore node is upserted. Both directions must fire —
|
||||
# the arrival order of chat messages vs contact advertisements is
|
||||
# not guaranteed and may differ across co-operating ingestors that
|
||||
# share this database. See issue #755.
|
||||
if protocol == "meshcore" && long_name && !long_name.empty?
|
||||
if synthetic == 0
|
||||
merge_synthetic_nodes(db, node_id, long_name)
|
||||
else
|
||||
merge_into_real_node(db, node_id, long_name)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
# Migrate messages from synthetic placeholder nodes to a newly confirmed
|
||||
# real node, then remove the placeholders.
|
||||
#
|
||||
# Called inside a transaction from +upsert_node+ when a real (non-synthetic)
|
||||
# MeshCore node with the same +long_name+ is upserted.
|
||||
#
|
||||
# Only +messages.from_id+ is migrated. Synthetic nodes are placeholders
|
||||
# created solely from parsed channel message sender names, so they cannot
|
||||
# have associated positions, telemetry, neighbors, or traces — those tables
|
||||
# are intentionally left untouched.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database connection.
|
||||
# @param real_node_id [String] canonical node ID for the real contact.
|
||||
# @param long_name [String] long name to match against synthetic rows.
|
||||
# @return [void]
|
||||
def merge_synthetic_nodes(db, real_node_id, long_name)
|
||||
# long_name is user-editable and not unique across pubkeys — two real
|
||||
# meshcore devices can legitimately share the same display name. When
|
||||
# that happens we cannot tell which real node a given chat-derived
|
||||
# synthetic was acting as placeholder for, so any merge would risk
|
||||
# mis-attributing messages. Bail out and leave the synthetic intact.
|
||||
other_real = db.execute(
|
||||
"SELECT 1 FROM nodes WHERE long_name = ? AND synthetic = 0 AND protocol = 'meshcore' AND node_id != ? LIMIT 1",
|
||||
[long_name, real_node_id],
|
||||
).first
|
||||
return if other_real
|
||||
|
||||
synthetic_ids = db.execute(
|
||||
"SELECT node_id FROM nodes WHERE long_name = ? AND synthetic = 1 AND protocol = 'meshcore' AND node_id != ?",
|
||||
[long_name, real_node_id],
|
||||
).map { |row| row[0] }
|
||||
|
||||
synthetic_ids.each do |synthetic_id|
|
||||
db.execute(
|
||||
"UPDATE messages SET from_id = ? WHERE from_id = ?",
|
||||
[real_node_id, synthetic_id],
|
||||
)
|
||||
db.execute(
|
||||
"DELETE FROM nodes WHERE node_id = ? AND synthetic = 1",
|
||||
[synthetic_id],
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
# Reverse of +merge_synthetic_nodes+: when a synthetic placeholder is
|
||||
# upserted for a MeshCore sender whose real contact advertisement has
|
||||
# already been stored (e.g. by a co-operating ingestor that saw the
|
||||
# advertisement first), migrate any messages from the synthetic id to the
|
||||
# real id and drop the synthetic row.
|
||||
#
|
||||
# Fixes duplication bug #755 where a chat-derived synthetic node and a
|
||||
# pubkey-derived real node coexisted because the forward merge only fired
|
||||
# on real-node upserts and never back-filled late-arriving synthetics.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database connection.
|
||||
# @param synthetic_node_id [String] canonical node ID of the synthetic placeholder being upserted.
|
||||
# @param long_name [String] long name to match against existing real rows.
|
||||
# @return [void]
|
||||
def merge_into_real_node(db, synthetic_node_id, long_name)
|
||||
# Index by [0] rather than the hash key so this works whether the db
|
||||
# handle was opened with results_as_hash = true or not.
|
||||
real_rows = db.execute(
|
||||
"SELECT node_id FROM nodes WHERE long_name = ? AND synthetic = 0 AND protocol = 'meshcore' AND node_id != ? LIMIT 2",
|
||||
[long_name, synthetic_node_id],
|
||||
)
|
||||
# Ambiguous name: two distinct real meshcore devices share this
|
||||
# long_name. The synthetic placeholder could legitimately represent
|
||||
# either, so we cannot pick one without risking mis-attribution. Leave
|
||||
# the synthetic in place; an operator can resolve the duplicate
|
||||
# manually.
|
||||
return if real_rows.length > 1
|
||||
|
||||
row = real_rows.first
|
||||
return unless row
|
||||
|
||||
real_node_id = row[0]
|
||||
return unless real_node_id
|
||||
|
||||
db.execute(
|
||||
"UPDATE messages SET from_id = ? WHERE from_id = ?",
|
||||
[real_node_id, synthetic_node_id],
|
||||
)
|
||||
db.execute(
|
||||
"DELETE FROM nodes WHERE node_id = ? AND synthetic = 1",
|
||||
[synthetic_node_id],
|
||||
)
|
||||
end
|
||||
|
||||
# Update node row columns from a freshly observed position record.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param node_id [String, nil] canonical node identifier.
|
||||
# @param node_num [Integer, nil] numeric node identifier.
|
||||
# @param rx_time [Integer, nil] receive time.
|
||||
# @param position_time [Integer, nil] timestamp from the position payload.
|
||||
# @param location_source [String, nil] +location_source+ enum value.
|
||||
# @param precision_bits [Integer, nil] horizontal precision bits.
|
||||
# @param latitude [Float, nil] decoded latitude.
|
||||
# @param longitude [Float, nil] decoded longitude.
|
||||
# @param altitude [Float, nil] decoded altitude.
|
||||
# @param snr [Float, nil] signal-to-noise ratio.
|
||||
# @return [void]
|
||||
def update_node_from_position(db, node_id, node_num, rx_time, position_time, location_source, precision_bits, latitude, longitude, altitude, snr)
|
||||
num = coerce_integer(node_num)
|
||||
id = string_or_nil(node_id)
|
||||
if id&.start_with?("!")
|
||||
id = "!#{id.delete_prefix("!").downcase}"
|
||||
end
|
||||
id ||= format("!%08x", num & 0xFFFFFFFF) if num
|
||||
return unless id
|
||||
|
||||
now = Time.now.to_i
|
||||
rx = coerce_integer(rx_time) || now
|
||||
rx = now if rx && rx > now
|
||||
pos_time = coerce_integer(position_time)
|
||||
pos_time = nil if pos_time && pos_time > now
|
||||
last_heard = [rx, pos_time].compact.max || rx
|
||||
last_heard = now if last_heard && last_heard > now
|
||||
|
||||
loc = string_or_nil(location_source)
|
||||
lat = coerce_float(latitude)
|
||||
lon = coerce_float(longitude)
|
||||
alt = coerce_float(altitude)
|
||||
precision = coerce_integer(precision_bits)
|
||||
snr_val = coerce_float(snr)
|
||||
|
||||
update_prometheus_metrics(node_id, nil, nil, nil, {
|
||||
"latitude" => lat,
|
||||
"longitude" => lon,
|
||||
"altitude" => alt,
|
||||
})
|
||||
|
||||
row = [
|
||||
id,
|
||||
num,
|
||||
last_heard,
|
||||
last_heard,
|
||||
pos_time,
|
||||
loc,
|
||||
precision,
|
||||
lat,
|
||||
lon,
|
||||
alt,
|
||||
snr_val,
|
||||
]
|
||||
with_busy_retry do
|
||||
db.execute <<~SQL, row
|
||||
INSERT INTO nodes(node_id,num,last_heard,first_heard,position_time,location_source,precision_bits,latitude,longitude,altitude,snr)
|
||||
VALUES (?,?,?,?,?,?,?,?,?,?,?)
|
||||
ON CONFLICT(node_id) DO UPDATE SET
|
||||
num=COALESCE(excluded.num,nodes.num),
|
||||
snr=COALESCE(excluded.snr,nodes.snr),
|
||||
last_heard=MAX(COALESCE(nodes.last_heard,0),COALESCE(excluded.last_heard,0)),
|
||||
first_heard=COALESCE(nodes.first_heard, excluded.first_heard, excluded.last_heard),
|
||||
position_time=CASE
|
||||
WHEN COALESCE(excluded.position_time,0) >= COALESCE(nodes.position_time,0)
|
||||
THEN excluded.position_time
|
||||
ELSE nodes.position_time
|
||||
END,
|
||||
location_source=CASE
|
||||
WHEN COALESCE(excluded.position_time,0) >= COALESCE(nodes.position_time,0)
|
||||
AND excluded.location_source IS NOT NULL
|
||||
THEN excluded.location_source
|
||||
ELSE nodes.location_source
|
||||
END,
|
||||
precision_bits=CASE
|
||||
WHEN COALESCE(excluded.position_time,0) >= COALESCE(nodes.position_time,0)
|
||||
AND excluded.precision_bits IS NOT NULL
|
||||
THEN excluded.precision_bits
|
||||
ELSE nodes.precision_bits
|
||||
END,
|
||||
latitude=CASE
|
||||
WHEN COALESCE(excluded.position_time,0) >= COALESCE(nodes.position_time,0)
|
||||
AND excluded.latitude IS NOT NULL
|
||||
THEN excluded.latitude
|
||||
ELSE nodes.latitude
|
||||
END,
|
||||
longitude=CASE
|
||||
WHEN COALESCE(excluded.position_time,0) >= COALESCE(nodes.position_time,0)
|
||||
AND excluded.longitude IS NOT NULL
|
||||
THEN excluded.longitude
|
||||
ELSE nodes.longitude
|
||||
END,
|
||||
altitude=CASE
|
||||
WHEN COALESCE(excluded.position_time,0) >= COALESCE(nodes.position_time,0)
|
||||
AND excluded.altitude IS NOT NULL
|
||||
THEN excluded.altitude
|
||||
ELSE nodes.altitude
|
||||
END
|
||||
SQL
|
||||
end
|
||||
end
|
||||
|
||||
# Update node columns based on metrics included in a telemetry packet.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param node_id [String, nil] canonical node identifier.
|
||||
# @param node_num [Integer, nil] numeric node identifier.
|
||||
# @param rx_time [Integer, nil] receive time used as +last_heard+.
|
||||
# @param metrics [Hash] decoded telemetry metric map.
|
||||
# @param lora_freq [Integer, nil] optional LoRa frequency.
|
||||
# @param modem_preset [String, nil] optional modem preset.
|
||||
# @param protocol [String] protocol identifier (default +meshtastic+).
|
||||
# @return [void]
|
||||
def update_node_from_telemetry(
|
||||
db,
|
||||
node_id,
|
||||
node_num,
|
||||
rx_time,
|
||||
metrics = {},
|
||||
lora_freq: nil,
|
||||
modem_preset: nil,
|
||||
protocol: "meshtastic"
|
||||
)
|
||||
num = coerce_integer(node_num)
|
||||
id = string_or_nil(node_id)
|
||||
if id&.start_with?("!")
|
||||
id = "!#{id.delete_prefix("!").downcase}"
|
||||
end
|
||||
id ||= format("!%08x", num & 0xFFFFFFFF) if num
|
||||
return unless id
|
||||
|
||||
ensure_unknown_node(db, id, num, heard_time: rx_time, protocol: protocol)
|
||||
touch_node_last_seen(
|
||||
db,
|
||||
id,
|
||||
num,
|
||||
rx_time: rx_time,
|
||||
source: :telemetry,
|
||||
lora_freq: lora_freq,
|
||||
modem_preset: modem_preset,
|
||||
)
|
||||
|
||||
battery = coerce_float(metrics[:battery_level] || metrics["battery_level"])
|
||||
voltage = coerce_float(metrics[:voltage] || metrics["voltage"])
|
||||
channel_util = coerce_float(metrics[:channel_utilization] || metrics["channel_utilization"])
|
||||
air_util_tx = coerce_float(metrics[:air_util_tx] || metrics["air_util_tx"])
|
||||
uptime = coerce_integer(metrics[:uptime_seconds] || metrics["uptime_seconds"])
|
||||
|
||||
update_prometheus_metrics(node_id, nil, nil, {
|
||||
"batteryLevel" => battery,
|
||||
"voltage" => voltage,
|
||||
"uptimeSeconds" => uptime,
|
||||
"channelUtilization" => channel_util,
|
||||
"airUtilTx" => air_util_tx,
|
||||
}, nil)
|
||||
|
||||
assignments = []
|
||||
params = []
|
||||
|
||||
if num
|
||||
assignments << "num = ?"
|
||||
params << num
|
||||
end
|
||||
|
||||
metric_updates = {
|
||||
"battery_level" => battery,
|
||||
"voltage" => voltage,
|
||||
"channel_utilization" => channel_util,
|
||||
"air_util_tx" => air_util_tx,
|
||||
"uptime_seconds" => uptime,
|
||||
}
|
||||
|
||||
metric_updates.each do |column, value|
|
||||
next if value.nil?
|
||||
|
||||
assignments << "#{column} = ?"
|
||||
params << value
|
||||
end
|
||||
|
||||
return if assignments.empty?
|
||||
|
||||
assignments_sql = assignments.join(", ")
|
||||
params << id
|
||||
|
||||
with_busy_retry do
|
||||
db.execute("UPDATE nodes SET #{assignments_sql} WHERE node_id = ?", params)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,226 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Persist a position payload, populate the +nodes+ table for newly seen
|
||||
# senders, and update node rows with the freshest GPS fields.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param payload [Hash] inbound position payload.
|
||||
# @param protocol_cache [Hash, nil] optional per-batch ingestor protocol cache.
|
||||
# @return [void]
|
||||
def insert_position(db, payload, protocol_cache: nil)
|
||||
pos_id = coerce_integer(payload["id"] || payload["packet_id"])
|
||||
return unless pos_id
|
||||
|
||||
now = Time.now.to_i
|
||||
rx_time = coerce_integer(payload["rx_time"])
|
||||
rx_time = now if rx_time.nil? || rx_time > now
|
||||
rx_iso = string_or_nil(payload["rx_iso"])
|
||||
rx_iso ||= Time.at(rx_time).utc.iso8601
|
||||
|
||||
raw_node_id = payload["node_id"] || payload["from_id"] || payload["from"]
|
||||
raw_node_num = coerce_integer(payload["node_num"]) || coerce_integer(payload["num"])
|
||||
|
||||
canonical_parts = canonical_node_parts(raw_node_id, raw_node_num)
|
||||
if canonical_parts
|
||||
node_id, node_num, = canonical_parts
|
||||
else
|
||||
node_id = string_or_nil(raw_node_id)
|
||||
node_id = "!#{node_id.delete_prefix("!").downcase}" if node_id&.start_with?("!")
|
||||
node_id ||= format("!%08x", raw_node_num & 0xFFFFFFFF) if node_id.nil? && raw_node_num
|
||||
|
||||
payload_for_num = payload.is_a?(Hash) ? payload.dup : {}
|
||||
payload_for_num["num"] ||= raw_node_num if raw_node_num
|
||||
node_num = resolve_node_num(node_id, payload_for_num)
|
||||
node_num ||= raw_node_num
|
||||
canonical = normalize_node_id(db, node_id || node_num)
|
||||
node_id = canonical if canonical
|
||||
end
|
||||
|
||||
lora_freq = coerce_integer(payload["lora_freq"] || payload["loraFrequency"])
|
||||
modem_preset = string_or_nil(payload["modem_preset"] || payload["modemPreset"])
|
||||
ingestor = string_or_nil(payload["ingestor"])
|
||||
protocol = resolve_protocol(db, ingestor, cache: protocol_cache)
|
||||
|
||||
ensure_unknown_node(db, node_id || node_num, node_num, heard_time: rx_time, protocol: protocol)
|
||||
touch_node_last_seen(
|
||||
db,
|
||||
node_id || node_num,
|
||||
node_num,
|
||||
rx_time: rx_time,
|
||||
source: :position,
|
||||
lora_freq: lora_freq,
|
||||
modem_preset: modem_preset,
|
||||
)
|
||||
|
||||
to_id = string_or_nil(payload["to_id"] || payload["to"])
|
||||
|
||||
position_section = payload["position"].is_a?(Hash) ? payload["position"] : {}
|
||||
|
||||
lat = coerce_float(payload["latitude"]) || coerce_float(position_section["latitude"])
|
||||
lon = coerce_float(payload["longitude"]) || coerce_float(position_section["longitude"])
|
||||
alt = coerce_float(payload["altitude"]) || coerce_float(position_section["altitude"])
|
||||
|
||||
lat ||= begin
|
||||
lat_i = coerce_integer(position_section["latitudeI"] || position_section["latitude_i"] || position_section.dig("raw", "latitude_i"))
|
||||
lat_i ? lat_i / 1e7 : nil
|
||||
end
|
||||
lon ||= begin
|
||||
lon_i = coerce_integer(position_section["longitudeI"] || position_section["longitude_i"] || position_section.dig("raw", "longitude_i"))
|
||||
lon_i ? lon_i / 1e7 : nil
|
||||
end
|
||||
alt ||= coerce_float(position_section.dig("raw", "altitude"))
|
||||
|
||||
position_time = coerce_integer(
|
||||
payload["position_time"] ||
|
||||
position_section["time"] ||
|
||||
position_section.dig("raw", "time"),
|
||||
)
|
||||
|
||||
location_source = string_or_nil(
|
||||
payload["location_source"] ||
|
||||
payload["locationSource"] ||
|
||||
position_section["location_source"] ||
|
||||
position_section["locationSource"] ||
|
||||
position_section.dig("raw", "location_source"),
|
||||
)
|
||||
|
||||
precision_bits = coerce_integer(
|
||||
payload["precision_bits"] ||
|
||||
payload["precisionBits"] ||
|
||||
position_section["precision_bits"] ||
|
||||
position_section["precisionBits"] ||
|
||||
position_section.dig("raw", "precision_bits"),
|
||||
)
|
||||
|
||||
sats_in_view = coerce_integer(
|
||||
payload["sats_in_view"] ||
|
||||
payload["satsInView"] ||
|
||||
position_section["sats_in_view"] ||
|
||||
position_section["satsInView"] ||
|
||||
position_section.dig("raw", "sats_in_view"),
|
||||
)
|
||||
|
||||
pdop = coerce_float(
|
||||
payload["pdop"] ||
|
||||
payload["PDOP"] ||
|
||||
position_section["pdop"] ||
|
||||
position_section["PDOP"] ||
|
||||
position_section.dig("raw", "PDOP") ||
|
||||
position_section.dig("raw", "pdop"),
|
||||
)
|
||||
|
||||
ground_speed = coerce_float(
|
||||
payload["ground_speed"] ||
|
||||
payload["groundSpeed"] ||
|
||||
position_section["ground_speed"] ||
|
||||
position_section["groundSpeed"] ||
|
||||
position_section.dig("raw", "ground_speed"),
|
||||
)
|
||||
|
||||
ground_track = coerce_float(
|
||||
payload["ground_track"] ||
|
||||
payload["groundTrack"] ||
|
||||
position_section["ground_track"] ||
|
||||
position_section["groundTrack"] ||
|
||||
position_section.dig("raw", "ground_track"),
|
||||
)
|
||||
|
||||
snr = coerce_float(payload["snr"] || payload["rx_snr"] || payload["rxSnr"])
|
||||
rssi = coerce_integer(payload["rssi"] || payload["rx_rssi"] || payload["rxRssi"])
|
||||
hop_limit = coerce_integer(payload["hop_limit"] || payload["hopLimit"])
|
||||
bitfield = coerce_integer(payload["bitfield"])
|
||||
|
||||
payload_b64 = string_or_nil(payload["payload_b64"] || payload["payload"])
|
||||
payload_b64 ||= string_or_nil(position_section.dig("payload", "__bytes_b64__"))
|
||||
|
||||
row = [
|
||||
pos_id,
|
||||
node_id,
|
||||
node_num,
|
||||
rx_time,
|
||||
rx_iso,
|
||||
position_time,
|
||||
to_id,
|
||||
lat,
|
||||
lon,
|
||||
alt,
|
||||
location_source,
|
||||
precision_bits,
|
||||
sats_in_view,
|
||||
pdop,
|
||||
ground_speed,
|
||||
ground_track,
|
||||
snr,
|
||||
rssi,
|
||||
hop_limit,
|
||||
bitfield,
|
||||
payload_b64,
|
||||
ingestor,
|
||||
protocol,
|
||||
]
|
||||
|
||||
with_busy_retry do
|
||||
db.execute <<~SQL, row
|
||||
INSERT INTO positions(id,node_id,node_num,rx_time,rx_iso,position_time,to_id,latitude,longitude,altitude,location_source,
|
||||
precision_bits,sats_in_view,pdop,ground_speed,ground_track,snr,rssi,hop_limit,bitfield,payload_b64,ingestor,protocol)
|
||||
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
|
||||
ON CONFLICT(id) DO UPDATE SET
|
||||
node_id=COALESCE(excluded.node_id,positions.node_id),
|
||||
node_num=COALESCE(excluded.node_num,positions.node_num),
|
||||
rx_time=excluded.rx_time,
|
||||
rx_iso=excluded.rx_iso,
|
||||
position_time=COALESCE(excluded.position_time,positions.position_time),
|
||||
to_id=COALESCE(excluded.to_id,positions.to_id),
|
||||
latitude=COALESCE(excluded.latitude,positions.latitude),
|
||||
longitude=COALESCE(excluded.longitude,positions.longitude),
|
||||
altitude=COALESCE(excluded.altitude,positions.altitude),
|
||||
location_source=COALESCE(excluded.location_source,positions.location_source),
|
||||
precision_bits=COALESCE(excluded.precision_bits,positions.precision_bits),
|
||||
sats_in_view=COALESCE(excluded.sats_in_view,positions.sats_in_view),
|
||||
pdop=COALESCE(excluded.pdop,positions.pdop),
|
||||
ground_speed=COALESCE(excluded.ground_speed,positions.ground_speed),
|
||||
ground_track=COALESCE(excluded.ground_track,positions.ground_track),
|
||||
snr=COALESCE(excluded.snr,positions.snr),
|
||||
rssi=COALESCE(excluded.rssi,positions.rssi),
|
||||
hop_limit=COALESCE(excluded.hop_limit,positions.hop_limit),
|
||||
bitfield=COALESCE(excluded.bitfield,positions.bitfield),
|
||||
payload_b64=COALESCE(excluded.payload_b64,positions.payload_b64),
|
||||
ingestor=COALESCE(NULLIF(positions.ingestor,''), excluded.ingestor),
|
||||
protocol=COALESCE(NULLIF(positions.protocol,'meshtastic'), excluded.protocol)
|
||||
SQL
|
||||
end
|
||||
|
||||
update_node_from_position(
|
||||
db,
|
||||
node_id,
|
||||
node_num,
|
||||
rx_time,
|
||||
position_time,
|
||||
location_source,
|
||||
precision_bits,
|
||||
lat,
|
||||
lon,
|
||||
alt,
|
||||
snr,
|
||||
)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,50 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Look up the protocol registered by a given ingestor node.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param ingestor_node_id [String, nil] the node_id of the reporting ingestor.
|
||||
# @param cache [Hash, nil] optional per-request memoization hash; pass a shared
|
||||
# Hash instance across a batch to avoid redundant DB lookups per record.
|
||||
# @return [String] protocol string; defaults to "meshtastic" when absent or unknown.
|
||||
def resolve_protocol(db, ingestor_node_id, cache: nil)
|
||||
return "meshtastic" if ingestor_node_id.nil? || ingestor_node_id.to_s.strip.empty?
|
||||
|
||||
if cache
|
||||
return cache[ingestor_node_id] if cache.key?(ingestor_node_id)
|
||||
|
||||
result = db.get_first_value(
|
||||
"SELECT protocol FROM ingestors WHERE node_id = ? LIMIT 1",
|
||||
[ingestor_node_id],
|
||||
) || "meshtastic"
|
||||
cache[ingestor_node_id] = result
|
||||
return result
|
||||
end
|
||||
|
||||
db.get_first_value(
|
||||
"SELECT protocol FROM ingestors WHERE node_id = ? LIMIT 1",
|
||||
[ingestor_node_id],
|
||||
) || "meshtastic"
|
||||
end
|
||||
|
||||
private :resolve_protocol
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,69 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Halt the current request with HTTP 403 unless the request carries a
|
||||
# bearer token that securely matches +API_TOKEN+.
|
||||
#
|
||||
# @return [void]
|
||||
def require_token!
|
||||
token = ENV["API_TOKEN"]
|
||||
provided = request.env["HTTP_AUTHORIZATION"].to_s.sub(/^Bearer\s+/i, "")
|
||||
halt 403, { error: "Forbidden" }.to_json unless token && !token.empty? && secure_token_match?(token, provided)
|
||||
end
|
||||
|
||||
# Constant-time comparison of two API tokens to mitigate timing attacks.
|
||||
#
|
||||
# @param expected [String] expected token from configuration.
|
||||
# @param provided [String] token supplied by the client.
|
||||
# @return [Boolean] true when the tokens match in constant time.
|
||||
def secure_token_match?(expected, provided)
|
||||
return false unless expected.is_a?(String) && provided.is_a?(String)
|
||||
|
||||
expected_bytes = expected.b
|
||||
provided_bytes = provided.b
|
||||
return false unless expected_bytes.bytesize == provided_bytes.bytesize
|
||||
Rack::Utils.secure_compare(expected_bytes, provided_bytes)
|
||||
rescue Rack::Utils::SecurityError
|
||||
false
|
||||
end
|
||||
|
||||
# Read the request body up to a configured byte ceiling and halt with HTTP
|
||||
# 413 when the payload exceeds the limit.
|
||||
#
|
||||
# @param limit [Integer, nil] optional override; falls back to
|
||||
# +PotatoMesh::Config.max_json_body_bytes+ when nil or non-positive.
|
||||
# @return [String] raw request body.
|
||||
def read_json_body(limit: nil)
|
||||
max_bytes = limit || PotatoMesh::Config.max_json_body_bytes
|
||||
max_bytes = max_bytes.to_i
|
||||
if max_bytes <= 0
|
||||
max_bytes = PotatoMesh::Config.max_json_body_bytes
|
||||
end
|
||||
|
||||
body = request.body.read(max_bytes + 1)
|
||||
body = "" if body.nil?
|
||||
halt 413, { error: "payload too large" }.to_json if body.bytesize > max_bytes
|
||||
|
||||
body
|
||||
ensure
|
||||
request.body.rewind if request.body.respond_to?(:rewind)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,547 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Ordered list of telemetry metric definitions consulted by
|
||||
# +insert_telemetry+. Each entry is a tuple of
|
||||
# +[column_name, coercion_type, key_map]+, where +key_map+ specifies the
|
||||
# candidate field names for each source layer. Hoisted out of the method
|
||||
# body to keep +insert_telemetry+ scannable; the data is otherwise
|
||||
# identical to the inline definitions used previously.
|
||||
TELEMETRY_METRIC_DEFINITIONS = [
|
||||
[
|
||||
"battery_level",
|
||||
:float,
|
||||
{
|
||||
payload: %w[battery_level batteryLevel],
|
||||
telemetry: %w[batteryLevel],
|
||||
device: %w[battery_level batteryLevel],
|
||||
environment: %w[battery_level batteryLevel],
|
||||
},
|
||||
],
|
||||
[
|
||||
"voltage",
|
||||
:float,
|
||||
{
|
||||
payload: %w[voltage],
|
||||
telemetry: %w[voltage],
|
||||
device: %w[voltage],
|
||||
environment: %w[voltage],
|
||||
},
|
||||
],
|
||||
[
|
||||
"channel_utilization",
|
||||
:float,
|
||||
{
|
||||
payload: %w[channel_utilization channelUtilization],
|
||||
telemetry: %w[channelUtilization],
|
||||
device: %w[channel_utilization channelUtilization],
|
||||
},
|
||||
],
|
||||
[
|
||||
"air_util_tx",
|
||||
:float,
|
||||
{
|
||||
payload: %w[air_util_tx airUtilTx],
|
||||
telemetry: %w[airUtilTx],
|
||||
device: %w[air_util_tx airUtilTx],
|
||||
},
|
||||
],
|
||||
[
|
||||
"uptime_seconds",
|
||||
:integer,
|
||||
{
|
||||
payload: %w[uptime_seconds uptimeSeconds],
|
||||
telemetry: %w[uptimeSeconds],
|
||||
device: %w[uptime_seconds uptimeSeconds],
|
||||
},
|
||||
],
|
||||
[
|
||||
"temperature",
|
||||
:float,
|
||||
{
|
||||
payload: %w[temperature temperatureC tempC],
|
||||
telemetry: %w[temperature temperatureC tempC],
|
||||
environment: %w[temperature temperatureC temperature_c tempC],
|
||||
},
|
||||
],
|
||||
[
|
||||
"relative_humidity",
|
||||
:float,
|
||||
{
|
||||
payload: %w[relative_humidity relativeHumidity humidity],
|
||||
telemetry: %w[relative_humidity relativeHumidity humidity],
|
||||
environment: %w[relative_humidity relativeHumidity humidity],
|
||||
},
|
||||
],
|
||||
[
|
||||
"barometric_pressure",
|
||||
:float,
|
||||
{
|
||||
payload: %w[barometric_pressure barometricPressure pressure],
|
||||
telemetry: %w[barometric_pressure barometricPressure pressure],
|
||||
environment: %w[barometric_pressure barometricPressure pressure],
|
||||
},
|
||||
],
|
||||
[
|
||||
"gas_resistance",
|
||||
:float,
|
||||
{
|
||||
payload: %w[gas_resistance gasResistance],
|
||||
telemetry: %w[gas_resistance gasResistance],
|
||||
environment: %w[gas_resistance gasResistance],
|
||||
},
|
||||
],
|
||||
[
|
||||
"current",
|
||||
:float,
|
||||
{
|
||||
payload: %w[current current_ma currentMa],
|
||||
telemetry: %w[current current_ma currentMa],
|
||||
device: %w[current current_ma currentMa],
|
||||
environment: %w[current],
|
||||
},
|
||||
],
|
||||
[
|
||||
"iaq",
|
||||
:integer,
|
||||
{
|
||||
payload: %w[iaq iaqIndex iaq_index],
|
||||
telemetry: %w[iaq iaqIndex iaq_index],
|
||||
environment: %w[iaq iaqIndex iaq_index],
|
||||
},
|
||||
],
|
||||
[
|
||||
"distance",
|
||||
:float,
|
||||
{
|
||||
payload: %w[distance range rangeMeters],
|
||||
telemetry: %w[distance range rangeMeters],
|
||||
environment: %w[distance range rangeMeters],
|
||||
},
|
||||
],
|
||||
[
|
||||
"lux",
|
||||
:float,
|
||||
{
|
||||
payload: %w[lux illuminance lightLux],
|
||||
telemetry: %w[lux illuminance lightLux],
|
||||
environment: %w[lux illuminance lightLux],
|
||||
},
|
||||
],
|
||||
[
|
||||
"white_lux",
|
||||
:float,
|
||||
{
|
||||
payload: %w[white_lux whiteLux],
|
||||
telemetry: %w[white_lux whiteLux],
|
||||
environment: %w[white_lux whiteLux],
|
||||
},
|
||||
],
|
||||
[
|
||||
"ir_lux",
|
||||
:float,
|
||||
{
|
||||
payload: %w[ir_lux irLux],
|
||||
telemetry: %w[ir_lux irLux],
|
||||
environment: %w[ir_lux irLux],
|
||||
},
|
||||
],
|
||||
[
|
||||
"uv_lux",
|
||||
:float,
|
||||
{
|
||||
payload: %w[uv_lux uvLux uvIndex],
|
||||
telemetry: %w[uv_lux uvLux uvIndex],
|
||||
environment: %w[uv_lux uvLux uvIndex],
|
||||
},
|
||||
],
|
||||
[
|
||||
"wind_direction",
|
||||
:integer,
|
||||
{
|
||||
payload: %w[wind_direction windDirection],
|
||||
telemetry: %w[wind_direction windDirection],
|
||||
environment: %w[wind_direction windDirection],
|
||||
},
|
||||
],
|
||||
[
|
||||
"wind_speed",
|
||||
:float,
|
||||
{
|
||||
payload: %w[wind_speed windSpeed windSpeedMps],
|
||||
telemetry: %w[wind_speed windSpeed windSpeedMps],
|
||||
environment: %w[wind_speed windSpeed windSpeedMps],
|
||||
},
|
||||
],
|
||||
[
|
||||
"weight",
|
||||
:float,
|
||||
{
|
||||
payload: %w[weight mass],
|
||||
telemetry: %w[weight mass],
|
||||
environment: %w[weight mass],
|
||||
},
|
||||
],
|
||||
[
|
||||
"wind_gust",
|
||||
:float,
|
||||
{
|
||||
payload: %w[wind_gust windGust],
|
||||
telemetry: %w[wind_gust windGust],
|
||||
environment: %w[wind_gust windGust],
|
||||
},
|
||||
],
|
||||
[
|
||||
"wind_lull",
|
||||
:float,
|
||||
{
|
||||
payload: %w[wind_lull windLull],
|
||||
telemetry: %w[wind_lull windLull],
|
||||
environment: %w[wind_lull windLull],
|
||||
},
|
||||
],
|
||||
[
|
||||
"radiation",
|
||||
:float,
|
||||
{
|
||||
payload: %w[radiation radiationLevel],
|
||||
telemetry: %w[radiation radiationLevel],
|
||||
environment: %w[radiation radiationLevel],
|
||||
},
|
||||
],
|
||||
[
|
||||
"rainfall_1h",
|
||||
:float,
|
||||
{
|
||||
payload: %w[rainfall_1h rainfall1h rainfallOneHour],
|
||||
telemetry: %w[rainfall_1h rainfall1h rainfallOneHour],
|
||||
environment: %w[rainfall_1h rainfall1h rainfallOneHour],
|
||||
},
|
||||
],
|
||||
[
|
||||
"rainfall_24h",
|
||||
:float,
|
||||
{
|
||||
payload: %w[rainfall_24h rainfall24h rainfallTwentyFourHour],
|
||||
telemetry: %w[rainfall_24h rainfall24h rainfallTwentyFourHour],
|
||||
environment: %w[rainfall_24h rainfall24h rainfallTwentyFourHour],
|
||||
},
|
||||
],
|
||||
[
|
||||
"soil_moisture",
|
||||
:integer,
|
||||
{
|
||||
payload: %w[soil_moisture soilMoisture],
|
||||
telemetry: %w[soil_moisture soilMoisture],
|
||||
environment: %w[soil_moisture soilMoisture],
|
||||
},
|
||||
],
|
||||
[
|
||||
"soil_temperature",
|
||||
:float,
|
||||
{
|
||||
payload: %w[soil_temperature soilTemperature],
|
||||
telemetry: %w[soil_temperature soilTemperature],
|
||||
environment: %w[soil_temperature soilTemperature],
|
||||
},
|
||||
],
|
||||
].freeze
|
||||
|
||||
# Resolve a telemetry metric from the provided data sources.
|
||||
#
|
||||
# @param key_map [Hash{Symbol=>Array<String>}] ordered mapping of source names to candidate keys.
|
||||
# @param sources [Hash{Symbol=>Hash}] data structures to search for metric values.
|
||||
# @param type [Symbol] coercion strategy, ``:float`` or ``:integer``.
|
||||
# @return [Numeric, nil] coerced metric value or nil when no candidates exist.
|
||||
def resolve_numeric_metric(key_map, sources, type)
|
||||
key_map.each do |source, keys|
|
||||
next if keys.nil? || keys.empty?
|
||||
|
||||
data = sources[source]
|
||||
next unless data.is_a?(Hash)
|
||||
|
||||
keys.each do |name|
|
||||
next if name.nil?
|
||||
|
||||
key = name.to_s
|
||||
value = if data.key?(key)
|
||||
data[key]
|
||||
else
|
||||
sym_key = key.to_sym
|
||||
data.key?(sym_key) ? data[sym_key] : nil
|
||||
end
|
||||
|
||||
next if value.nil?
|
||||
|
||||
coerced = case type
|
||||
when :float
|
||||
coerce_float(value)
|
||||
when :integer
|
||||
coerce_integer(value)
|
||||
else
|
||||
value
|
||||
end
|
||||
|
||||
return coerced unless coerced.nil?
|
||||
end
|
||||
end
|
||||
|
||||
nil
|
||||
end
|
||||
|
||||
private :resolve_numeric_metric
|
||||
|
||||
# Persist a telemetry packet and refresh the related node row.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param payload [Hash] inbound telemetry payload.
|
||||
# @param protocol_cache [Hash, nil] optional per-batch ingestor protocol cache.
|
||||
# @return [void]
|
||||
def insert_telemetry(db, payload, protocol_cache: nil)
|
||||
return unless payload.is_a?(Hash)
|
||||
|
||||
telemetry_id = coerce_integer(payload["id"] || payload["packet_id"])
|
||||
return unless telemetry_id
|
||||
|
||||
now = Time.now.to_i
|
||||
rx_time = coerce_integer(payload["rx_time"])
|
||||
rx_time = now if rx_time.nil? || rx_time > now
|
||||
rx_iso = string_or_nil(payload["rx_iso"])
|
||||
rx_iso ||= Time.at(rx_time).utc.iso8601
|
||||
|
||||
raw_node_id = payload["node_id"] || payload["from_id"] || payload["from"]
|
||||
raw_node_num = coerce_integer(payload["node_num"]) || coerce_integer(payload["num"])
|
||||
|
||||
canonical_parts = canonical_node_parts(raw_node_id, raw_node_num)
|
||||
if canonical_parts
|
||||
node_id, node_num, = canonical_parts
|
||||
else
|
||||
node_id = string_or_nil(raw_node_id)
|
||||
node_id = "!#{node_id.delete_prefix("!").downcase}" if node_id&.start_with?("!")
|
||||
|
||||
payload_for_num = payload.dup
|
||||
payload_for_num["num"] ||= raw_node_num if raw_node_num
|
||||
node_num = resolve_node_num(node_id, payload_for_num)
|
||||
node_num ||= raw_node_num
|
||||
|
||||
canonical = normalize_node_id(db, node_id || node_num)
|
||||
node_id = canonical if canonical
|
||||
end
|
||||
|
||||
from_id = string_or_nil(payload["from_id"]) || node_id
|
||||
to_id = string_or_nil(payload["to_id"] || payload["to"])
|
||||
|
||||
telemetry_time = coerce_integer(payload["telemetry_time"] || payload["time"] || payload.dig("telemetry", "time"))
|
||||
telemetry_time = nil if telemetry_time && telemetry_time > now
|
||||
|
||||
channel = coerce_integer(payload["channel"])
|
||||
portnum = string_or_nil(payload["portnum"])
|
||||
hop_limit = coerce_integer(payload["hop_limit"] || payload["hopLimit"])
|
||||
snr = coerce_float(payload["snr"])
|
||||
rssi = coerce_integer(payload["rssi"])
|
||||
bitfield = coerce_integer(payload["bitfield"])
|
||||
payload_b64 = string_or_nil(payload["payload_b64"] || payload["payload"])
|
||||
lora_freq = coerce_integer(payload["lora_freq"] || payload["loraFrequency"])
|
||||
modem_preset = string_or_nil(payload["modem_preset"] || payload["modemPreset"])
|
||||
ingestor = string_or_nil(payload["ingestor"])
|
||||
protocol = resolve_protocol(db, ingestor, cache: protocol_cache)
|
||||
|
||||
telemetry_section = normalize_json_object(payload["telemetry"])
|
||||
device_metrics = normalize_json_object(payload["device_metrics"] || payload["deviceMetrics"])
|
||||
device_metrics ||= normalize_json_object(telemetry_section["deviceMetrics"]) if telemetry_section&.key?("deviceMetrics")
|
||||
environment_metrics = normalize_json_object(payload["environment_metrics"] || payload["environmentMetrics"])
|
||||
environment_metrics ||= normalize_json_object(telemetry_section["environmentMetrics"]) if telemetry_section&.key?("environmentMetrics")
|
||||
power_metrics = normalize_json_object(payload["power_metrics"] || payload["powerMetrics"])
|
||||
power_metrics ||= normalize_json_object(telemetry_section["powerMetrics"]) if telemetry_section&.key?("powerMetrics")
|
||||
air_quality_metrics = normalize_json_object(payload["air_quality_metrics"] || payload["airQualityMetrics"])
|
||||
air_quality_metrics ||= normalize_json_object(telemetry_section["airQualityMetrics"]) if telemetry_section&.key?("airQualityMetrics")
|
||||
|
||||
telemetry_type = string_or_nil(payload["telemetry_type"])
|
||||
telemetry_type = nil unless VALID_TELEMETRY_TYPES.include?(telemetry_type)
|
||||
telemetry_type ||= if device_metrics&.any?
|
||||
"device"
|
||||
elsif environment_metrics&.any?
|
||||
"environment"
|
||||
elsif power_metrics&.any?
|
||||
"power"
|
||||
elsif air_quality_metrics&.any?
|
||||
"air_quality"
|
||||
end
|
||||
|
||||
sources = {
|
||||
payload: payload,
|
||||
telemetry: telemetry_section,
|
||||
device: device_metrics,
|
||||
environment: environment_metrics,
|
||||
}
|
||||
|
||||
metric_values = {}
|
||||
TELEMETRY_METRIC_DEFINITIONS.each do |column, type, key_map|
|
||||
value = resolve_numeric_metric(key_map, sources, type)
|
||||
metric_values[column] = value unless value.nil?
|
||||
end
|
||||
|
||||
battery_level = metric_values["battery_level"]
|
||||
voltage = metric_values["voltage"]
|
||||
channel_utilization = metric_values["channel_utilization"]
|
||||
air_util_tx = metric_values["air_util_tx"]
|
||||
uptime_seconds = metric_values["uptime_seconds"]
|
||||
temperature = metric_values["temperature"]
|
||||
relative_humidity = metric_values["relative_humidity"]
|
||||
barometric_pressure = metric_values["barometric_pressure"]
|
||||
gas_resistance = metric_values["gas_resistance"]
|
||||
current = metric_values["current"]
|
||||
iaq = metric_values["iaq"]
|
||||
distance = metric_values["distance"]
|
||||
lux = metric_values["lux"]
|
||||
white_lux = metric_values["white_lux"]
|
||||
ir_lux = metric_values["ir_lux"]
|
||||
uv_lux = metric_values["uv_lux"]
|
||||
wind_direction = metric_values["wind_direction"]
|
||||
wind_speed = metric_values["wind_speed"]
|
||||
weight = metric_values["weight"]
|
||||
wind_gust = metric_values["wind_gust"]
|
||||
wind_lull = metric_values["wind_lull"]
|
||||
radiation = metric_values["radiation"]
|
||||
rainfall_1h = metric_values["rainfall_1h"]
|
||||
rainfall_24h = metric_values["rainfall_24h"]
|
||||
soil_moisture = metric_values["soil_moisture"]
|
||||
soil_temperature = metric_values["soil_temperature"]
|
||||
|
||||
row = [
|
||||
telemetry_id,
|
||||
node_id,
|
||||
node_num,
|
||||
from_id,
|
||||
to_id,
|
||||
rx_time,
|
||||
rx_iso,
|
||||
telemetry_time,
|
||||
channel,
|
||||
portnum,
|
||||
hop_limit,
|
||||
snr,
|
||||
rssi,
|
||||
bitfield,
|
||||
payload_b64,
|
||||
battery_level,
|
||||
voltage,
|
||||
channel_utilization,
|
||||
air_util_tx,
|
||||
uptime_seconds,
|
||||
temperature,
|
||||
relative_humidity,
|
||||
barometric_pressure,
|
||||
gas_resistance,
|
||||
current,
|
||||
iaq,
|
||||
distance,
|
||||
lux,
|
||||
white_lux,
|
||||
ir_lux,
|
||||
uv_lux,
|
||||
wind_direction,
|
||||
wind_speed,
|
||||
weight,
|
||||
wind_gust,
|
||||
wind_lull,
|
||||
radiation,
|
||||
rainfall_1h,
|
||||
rainfall_24h,
|
||||
soil_moisture,
|
||||
soil_temperature,
|
||||
ingestor,
|
||||
protocol,
|
||||
telemetry_type,
|
||||
]
|
||||
|
||||
placeholders = Array.new(row.length, "?").join(",")
|
||||
|
||||
with_busy_retry do
|
||||
db.execute <<~SQL, row
|
||||
INSERT INTO telemetry(id,node_id,node_num,from_id,to_id,rx_time,rx_iso,telemetry_time,channel,portnum,hop_limit,snr,rssi,bitfield,payload_b64,
|
||||
battery_level,voltage,channel_utilization,air_util_tx,uptime_seconds,temperature,relative_humidity,barometric_pressure,gas_resistance,current,iaq,distance,lux,white_lux,ir_lux,uv_lux,wind_direction,wind_speed,weight,wind_gust,wind_lull,radiation,rainfall_1h,rainfall_24h,soil_moisture,soil_temperature,ingestor,protocol,telemetry_type)
|
||||
VALUES (#{placeholders})
|
||||
ON CONFLICT(id) DO UPDATE SET
|
||||
node_id=COALESCE(excluded.node_id,telemetry.node_id),
|
||||
node_num=COALESCE(excluded.node_num,telemetry.node_num),
|
||||
from_id=COALESCE(excluded.from_id,telemetry.from_id),
|
||||
to_id=COALESCE(excluded.to_id,telemetry.to_id),
|
||||
rx_time=excluded.rx_time,
|
||||
rx_iso=excluded.rx_iso,
|
||||
telemetry_time=COALESCE(excluded.telemetry_time,telemetry.telemetry_time),
|
||||
channel=COALESCE(excluded.channel,telemetry.channel),
|
||||
portnum=COALESCE(excluded.portnum,telemetry.portnum),
|
||||
hop_limit=COALESCE(excluded.hop_limit,telemetry.hop_limit),
|
||||
snr=COALESCE(excluded.snr,telemetry.snr),
|
||||
rssi=COALESCE(excluded.rssi,telemetry.rssi),
|
||||
bitfield=COALESCE(excluded.bitfield,telemetry.bitfield),
|
||||
payload_b64=COALESCE(excluded.payload_b64,telemetry.payload_b64),
|
||||
battery_level=COALESCE(excluded.battery_level,telemetry.battery_level),
|
||||
voltage=COALESCE(excluded.voltage,telemetry.voltage),
|
||||
channel_utilization=COALESCE(excluded.channel_utilization,telemetry.channel_utilization),
|
||||
air_util_tx=COALESCE(excluded.air_util_tx,telemetry.air_util_tx),
|
||||
uptime_seconds=COALESCE(excluded.uptime_seconds,telemetry.uptime_seconds),
|
||||
temperature=COALESCE(excluded.temperature,telemetry.temperature),
|
||||
relative_humidity=COALESCE(excluded.relative_humidity,telemetry.relative_humidity),
|
||||
barometric_pressure=COALESCE(excluded.barometric_pressure,telemetry.barometric_pressure),
|
||||
gas_resistance=COALESCE(excluded.gas_resistance,telemetry.gas_resistance),
|
||||
current=COALESCE(excluded.current,telemetry.current),
|
||||
iaq=COALESCE(excluded.iaq,telemetry.iaq),
|
||||
distance=COALESCE(excluded.distance,telemetry.distance),
|
||||
lux=COALESCE(excluded.lux,telemetry.lux),
|
||||
white_lux=COALESCE(excluded.white_lux,telemetry.white_lux),
|
||||
ir_lux=COALESCE(excluded.ir_lux,telemetry.ir_lux),
|
||||
uv_lux=COALESCE(excluded.uv_lux,telemetry.uv_lux),
|
||||
wind_direction=COALESCE(excluded.wind_direction,telemetry.wind_direction),
|
||||
wind_speed=COALESCE(excluded.wind_speed,telemetry.wind_speed),
|
||||
weight=COALESCE(excluded.weight,telemetry.weight),
|
||||
wind_gust=COALESCE(excluded.wind_gust,telemetry.wind_gust),
|
||||
wind_lull=COALESCE(excluded.wind_lull,telemetry.wind_lull),
|
||||
radiation=COALESCE(excluded.radiation,telemetry.radiation),
|
||||
rainfall_1h=COALESCE(excluded.rainfall_1h,telemetry.rainfall_1h),
|
||||
rainfall_24h=COALESCE(excluded.rainfall_24h,telemetry.rainfall_24h),
|
||||
soil_moisture=COALESCE(excluded.soil_moisture,telemetry.soil_moisture),
|
||||
soil_temperature=COALESCE(excluded.soil_temperature,telemetry.soil_temperature),
|
||||
ingestor=COALESCE(NULLIF(telemetry.ingestor,''), excluded.ingestor),
|
||||
protocol=COALESCE(NULLIF(telemetry.protocol,'meshtastic'), excluded.protocol),
|
||||
telemetry_type=COALESCE(excluded.telemetry_type,telemetry.telemetry_type)
|
||||
SQL
|
||||
end
|
||||
|
||||
update_node_from_telemetry(
|
||||
db,
|
||||
node_id,
|
||||
node_num,
|
||||
rx_time,
|
||||
{
|
||||
battery_level: battery_level,
|
||||
voltage: voltage,
|
||||
channel_utilization: channel_utilization,
|
||||
air_util_tx: air_util_tx,
|
||||
uptime_seconds: uptime_seconds,
|
||||
},
|
||||
lora_freq: lora_freq,
|
||||
modem_preset: modem_preset,
|
||||
protocol: protocol,
|
||||
)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,130 @@
|
||||
# Copyright © 2025-26 l5yth & contributors
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# frozen_string_literal: true
|
||||
|
||||
module PotatoMesh
|
||||
module App
|
||||
module DataProcessing
|
||||
# Normalise a traceroute hop entry to a numeric node identifier.
|
||||
#
|
||||
# @param hop [Object] raw hop entry from the payload.
|
||||
# @return [Integer, nil] coerced node ID or nil when the value is unusable.
|
||||
def coerce_trace_node_id(hop)
|
||||
case hop
|
||||
when Integer
|
||||
return hop
|
||||
when Numeric
|
||||
return hop.to_i
|
||||
when String
|
||||
trimmed = hop.strip
|
||||
return nil if trimmed.empty?
|
||||
return Integer(trimmed, 10) if trimmed.match?(/\A-?\d+\z/)
|
||||
|
||||
parts = canonical_node_parts(trimmed)
|
||||
return parts[1] if parts
|
||||
when Hash
|
||||
candidate = hop["node_id"] || hop[:node_id] || hop["id"] || hop[:id] || hop["num"] || hop[:num]
|
||||
return coerce_trace_node_id(candidate)
|
||||
end
|
||||
|
||||
nil
|
||||
end
|
||||
|
||||
# Extract hop identifiers from a traceroute payload preserving order.
|
||||
#
|
||||
# @param hops_value [Object] raw hops array or path collection.
|
||||
# @return [Array<Integer>] ordered list of coerced hop identifiers.
|
||||
def normalize_trace_hops(hops_value)
|
||||
return [] if hops_value.nil?
|
||||
|
||||
hop_entries = hops_value.is_a?(Array) ? hops_value : [hops_value]
|
||||
hop_entries.filter_map { |entry| coerce_trace_node_id(entry) }
|
||||
end
|
||||
|
||||
# Persist a traceroute observation and its hop path.
|
||||
#
|
||||
# @param db [SQLite3::Database] open database handle.
|
||||
# @param payload [Hash] traceroute payload as produced by the ingestor.
|
||||
# @param protocol_cache [Hash, nil] optional per-batch ingestor protocol cache.
|
||||
# @return [void]
|
||||
def insert_trace(db, payload, protocol_cache: nil)
|
||||
return unless payload.is_a?(Hash)
|
||||
|
||||
trace_identifier = coerce_integer(payload["id"] || payload["packet_id"] || payload["packetId"])
|
||||
trace_identifier ||= coerce_integer(payload["trace_id"])
|
||||
request_id = coerce_integer(payload["request_id"] || payload["req"])
|
||||
trace_identifier ||= request_id
|
||||
|
||||
now = Time.now.to_i
|
||||
rx_time = coerce_integer(payload["rx_time"])
|
||||
rx_time = now if rx_time.nil? || rx_time > now
|
||||
rx_iso = string_or_nil(payload["rx_iso"]) || Time.at(rx_time).utc.iso8601
|
||||
|
||||
metrics = normalize_json_object(payload["metrics"]) || {}
|
||||
src = coerce_integer(payload["src"] || payload["source"] || payload["from"])
|
||||
dest = coerce_integer(payload["dest"] || payload["destination"] || payload["to"])
|
||||
rssi = coerce_integer(payload["rssi"]) || coerce_integer(metrics["rssi"])
|
||||
snr = coerce_float(payload["snr"]) || coerce_float(metrics["snr"])
|
||||
elapsed_ms = coerce_integer(
|
||||
payload["elapsed_ms"] ||
|
||||
payload["latency_ms"] ||
|
||||
metrics&.[]("elapsed_ms") ||
|
||||
metrics&.[]("latency_ms") ||
|
||||
metrics&.[]("latencyMs"),
|
||||
)
|
||||
ingestor = string_or_nil(payload["ingestor"])
|
||||
protocol = resolve_protocol(db, ingestor, cache: protocol_cache)
|
||||
|
||||
hops_value = payload.key?("hops") ? payload["hops"] : payload["path"]
|
||||
hops = normalize_trace_hops(hops_value)
|
||||
|
||||
all_nodes = [src, dest, *hops].compact.uniq
|
||||
all_nodes.each do |node|
|
||||
ensure_unknown_node(db, node, node, heard_time: rx_time, protocol: protocol)
|
||||
touch_node_last_seen(db, node, node, rx_time: rx_time, source: :trace)
|
||||
end
|
||||
|
||||
with_busy_retry do
|
||||
db.execute <<~SQL, [trace_identifier, request_id, src, dest, rx_time, rx_iso, rssi, snr, elapsed_ms, ingestor, protocol]
|
||||
INSERT INTO traces(id, request_id, src, dest, rx_time, rx_iso, rssi, snr, elapsed_ms, ingestor, protocol)
|
||||
VALUES(?,?,?,?,?,?,?,?,?,?,?)
|
||||
ON CONFLICT(id) DO UPDATE SET
|
||||
request_id=COALESCE(excluded.request_id,traces.request_id),
|
||||
src=COALESCE(excluded.src,traces.src),
|
||||
dest=COALESCE(excluded.dest,traces.dest),
|
||||
rx_time=excluded.rx_time,
|
||||
rx_iso=excluded.rx_iso,
|
||||
rssi=COALESCE(excluded.rssi,traces.rssi),
|
||||
snr=COALESCE(excluded.snr,traces.snr),
|
||||
elapsed_ms=COALESCE(excluded.elapsed_ms,traces.elapsed_ms),
|
||||
ingestor=COALESCE(NULLIF(traces.ingestor,''), excluded.ingestor),
|
||||
protocol=COALESCE(NULLIF(traces.protocol,'meshtastic'), excluded.protocol)
|
||||
SQL
|
||||
|
||||
trace_id = trace_identifier || db.last_insert_row_id
|
||||
return unless trace_id
|
||||
|
||||
db.execute("DELETE FROM trace_hops WHERE trace_id = ?", [trace_id])
|
||||
hops.each_with_index do |hop_id, index|
|
||||
db.execute(
|
||||
"INSERT INTO trace_hops(trace_id, hop_index, node_id) VALUES(?,?,?)",
|
||||
[trace_id, index, hop_id],
|
||||
)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -1026,4 +1026,600 @@ RSpec.describe PotatoMesh::App::DataProcessing do
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Coverage gap-fillers for the post-split data_processing/ submodules.
|
||||
#
|
||||
# The PR that split data_processing.rb into focused submodules surfaced a
|
||||
# set of pre-existing untested branches as "uncovered patch lines". The
|
||||
# describes below pin behaviour for each of those branches so future
|
||||
# changes cannot silently regress them. Grouped by submodule for
|
||||
# discoverability.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# identity.rb — numeric and bare-hex paths in canonical_node_parts
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#canonical_node_parts (numeric and bare-hex paths)" do
|
||||
it "coerces a Float node_ref to its integer counterpart" do
|
||||
parts = dp.canonical_node_parts(42.7)
|
||||
expect(parts).not_to be_nil
|
||||
expect(parts[1]).to eq(42)
|
||||
end
|
||||
|
||||
it "parses a bare lowercase hex string without the ! sigil" do
|
||||
parts = dp.canonical_node_parts("aabbccdd")
|
||||
expect(parts).not_to be_nil
|
||||
expect(parts[0]).to eq("!aabbccdd")
|
||||
expect(parts[1]).to eq(0xaabbccdd)
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# traces.rb — coerce_trace_node_id type handling
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#coerce_trace_node_id" do
|
||||
it "coerces a Float hop to an integer" do
|
||||
expect(dp.coerce_trace_node_id(42.7)).to eq(42)
|
||||
end
|
||||
|
||||
it "returns nil for unsupported hop types" do
|
||||
expect(dp.coerce_trace_node_id([1, 2])).to be_nil
|
||||
expect(dp.coerce_trace_node_id(true)).to be_nil
|
||||
end
|
||||
|
||||
it "extracts node ids from hash hops" do
|
||||
expect(dp.coerce_trace_node_id({ "node_id" => 12345 })).to eq(12345)
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# telemetry.rb — resolve_numeric_metric default coercion + power fallback
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#resolve_numeric_metric (private)" do
|
||||
it "passes through values unchanged for unknown coercion types" do
|
||||
sources = { payload: { "raw" => "literal" } }
|
||||
key_map = { payload: %w[raw] }
|
||||
result = dp.send(:resolve_numeric_metric, key_map, sources, :raw)
|
||||
expect(result).to eq("literal")
|
||||
end
|
||||
end
|
||||
|
||||
describe "#insert_telemetry — telemetry_type fallback chain" do
|
||||
include_context "with isolated db"
|
||||
|
||||
it "tags the row as 'power' when only power_metrics are supplied" do
|
||||
db = open_db
|
||||
dp.upsert_node(db, "!aabbccdd", { "num" => 0xaabbccdd })
|
||||
dp.insert_telemetry(db, {
|
||||
"id" => 5005,
|
||||
"node_id" => "!aabbccdd",
|
||||
"rx_time" => now,
|
||||
"power_metrics" => { "ch1Voltage" => 3.7 },
|
||||
})
|
||||
expect(db.get_first_value("SELECT telemetry_type FROM telemetry WHERE id = 5005")).to eq("power")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# node_writes.rb — touch_node_last_seen falls back to fallback_num when
|
||||
# node_ref strips down to nothing.
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#touch_node_last_seen — fallback_num when node_ref is blank" do
|
||||
include_context "with isolated db"
|
||||
|
||||
it "resolves node_id from fallback_num and refreshes last_heard" do
|
||||
db = open_db
|
||||
dp.upsert_node(db, "!00003039", { "num" => 12345, "lastHeard" => now - 100 })
|
||||
dp.touch_node_last_seen(db, "", 12345, rx_time: now, source: :test)
|
||||
expect(
|
||||
db.get_first_value("SELECT last_heard FROM nodes WHERE node_id = '!00003039'"),
|
||||
).to eq(now)
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# ingestors.rb — SQLite3::SQLException rescue
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#upsert_ingestor — SQL exception handling" do
|
||||
include_context "with isolated db"
|
||||
|
||||
it "returns false when the upsert raises SQLite3::SQLException" do
|
||||
db = open_db
|
||||
allow(db).to receive(:execute).and_raise(SQLite3::SQLException.new("boom"))
|
||||
expect(
|
||||
dp.upsert_ingestor(db, {
|
||||
"node_id" => "!aabbccdd",
|
||||
"version" => "1.0",
|
||||
"start_time" => now,
|
||||
}),
|
||||
).to be(false)
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# request_helpers.rb — read_json_body resets to the configured cap when the
|
||||
# caller-supplied limit collapses to zero or a negative value.
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#read_json_body — limit fallback" do
|
||||
let(:body_harness_class) do
|
||||
Class.new do
|
||||
include PotatoMesh::App::DataProcessing
|
||||
include PotatoMesh::App::Helpers
|
||||
attr_accessor :request
|
||||
|
||||
def halt(*args)
|
||||
raise "unexpected halt: #{args.inspect}"
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
it "falls back to the configured cap when limit is non-positive" do
|
||||
body_text = "small body"
|
||||
fake_request = Struct.new(:body).new(StringIO.new(body_text))
|
||||
instance = body_harness_class.new
|
||||
instance.request = fake_request
|
||||
expect(instance.read_json_body(limit: 0)).to eq(body_text)
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# decrypted_payloads.rb — store_decrypted_payload returns false when the
|
||||
# decoder reports an unrecognised payload type.
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#store_decrypted_payload — unrecognised type" do
|
||||
include_context "with isolated db"
|
||||
|
||||
it "returns false for a decoded type the case statement does not handle" do
|
||||
db = open_db
|
||||
allow(PotatoMesh::App::Meshtastic::PayloadDecoder).to receive(:decode)
|
||||
.and_return({ "type" => "UNRECOGNIZED", "payload" => {} })
|
||||
decrypted = { payload: "\x00\x01".b, portnum: 3 }
|
||||
result = dp.store_decrypted_payload(
|
||||
db, {}, 555, decrypted,
|
||||
rx_time: now,
|
||||
rx_iso: Time.at(now).utc.iso8601,
|
||||
from_id: "!aabbccdd",
|
||||
to_id: "^all",
|
||||
channel: 0,
|
||||
portnum: 3,
|
||||
hop_limit: 5,
|
||||
snr: 1.0,
|
||||
rssi: -50,
|
||||
)
|
||||
expect(result).to be(false)
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# neighbors.rb — non-canonical reporter and neighbour entry resolution.
|
||||
# Exercises the else branch of canonical_node_parts at the top of
|
||||
# +insert_neighbors+ and the equivalent inside the per-neighbour loop.
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#insert_neighbors — non-canonical resolution paths" do
|
||||
include_context "with isolated db"
|
||||
|
||||
let(:neighbor_harness) do
|
||||
Class.new do
|
||||
include PotatoMesh::App::DataProcessing
|
||||
include PotatoMesh::App::Helpers
|
||||
|
||||
def debug_log(*); end
|
||||
|
||||
def warn_log(*); end
|
||||
|
||||
def with_busy_retry
|
||||
yield
|
||||
end
|
||||
|
||||
def update_prometheus_metrics(*); end
|
||||
|
||||
def prom_report_ids
|
||||
[]
|
||||
end
|
||||
|
||||
def private_mode?
|
||||
false
|
||||
end
|
||||
|
||||
def resolve_protocol(*)
|
||||
"meshtastic"
|
||||
end
|
||||
|
||||
def normalize_node_id(_db, ref)
|
||||
parts = canonical_node_parts(ref)
|
||||
parts ? parts[0] : nil
|
||||
end
|
||||
|
||||
def ensure_unknown_node(*); end
|
||||
|
||||
def touch_node_last_seen(*); end
|
||||
end.new
|
||||
end
|
||||
|
||||
it "resolves node_num from fallback when node_id strips to empty" do
|
||||
db = open_db
|
||||
neighbor_harness.insert_neighbors(db, {
|
||||
"node_id" => "",
|
||||
"node_num" => 12345,
|
||||
"rx_time" => now,
|
||||
"neighbors" => [],
|
||||
})
|
||||
# Empty neighbor list deletes any existing rows for the reporter.
|
||||
expect(db.get_first_value("SELECT COUNT(*) FROM neighbors")).to eq(0)
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
|
||||
it "tolerates a malformed !-prefixed node id by zeroing the node_num" do
|
||||
db = open_db
|
||||
neighbor_harness.insert_neighbors(db, {
|
||||
"node_id" => "!ZZZ",
|
||||
"rx_time" => now,
|
||||
"neighbors" => [],
|
||||
})
|
||||
expect(db.get_first_value("SELECT COUNT(*) FROM neighbors")).to eq(0)
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
|
||||
it "resolves a neighbour entry from neighbor_num when neighbor_id is blank" do
|
||||
db = open_db
|
||||
neighbor_harness.insert_neighbors(db, {
|
||||
"node_id" => "!aabbccdd",
|
||||
"rx_time" => now,
|
||||
"neighbors" => [
|
||||
{ "neighbor_id" => "", "neighbor_num" => 6789, "snr" => -3.0 },
|
||||
],
|
||||
})
|
||||
stored_id = db.get_first_value(
|
||||
"SELECT neighbor_id FROM neighbors WHERE node_id = '!aabbccdd'",
|
||||
)
|
||||
expect(stored_id).to start_with("!")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
|
||||
it "tolerates a malformed !-prefixed neighbour id without inserting it" do
|
||||
db = open_db
|
||||
neighbor_harness.insert_neighbors(db, {
|
||||
"node_id" => "!aabbccdd",
|
||||
"rx_time" => now,
|
||||
"neighbors" => [
|
||||
{ "neighbor_id" => "!ZZZ", "snr" => 5.0 },
|
||||
],
|
||||
})
|
||||
# The malformed neighbour normalizes to "!zzz" (lowercased) which is
|
||||
# still stored under that bogus id; the contract under test is only that
|
||||
# the function does not raise. Coverage of the rescue/else branches is
|
||||
# the goal.
|
||||
expect(db).not_to be_nil
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# messages.rb — canonical sender/recipient overrides and the rare
|
||||
# ConstraintException recovery path inside +insert_message+.
|
||||
# ---------------------------------------------------------------------------
|
||||
describe "#insert_message — canonical sender/recipient overrides" do
|
||||
include_context "with isolated db"
|
||||
|
||||
let(:message_harness) do
|
||||
Class.new do
|
||||
include PotatoMesh::App::DataProcessing
|
||||
include PotatoMesh::App::Helpers
|
||||
|
||||
def debug_log(*); end
|
||||
|
||||
def warn_log(*); end
|
||||
|
||||
def with_busy_retry
|
||||
yield
|
||||
end
|
||||
|
||||
def update_prometheus_metrics(*); end
|
||||
|
||||
def prom_report_ids
|
||||
[]
|
||||
end
|
||||
|
||||
def private_mode?
|
||||
false
|
||||
end
|
||||
|
||||
def resolve_protocol(*)
|
||||
"meshtastic"
|
||||
end
|
||||
|
||||
def normalize_node_id(_db, ref)
|
||||
parts = canonical_node_parts(ref)
|
||||
parts ? parts[0] : nil
|
||||
end
|
||||
|
||||
def ensure_unknown_node(*); end
|
||||
|
||||
def touch_node_last_seen(*); end
|
||||
end.new
|
||||
end
|
||||
|
||||
it "uses canonical_from_id when the raw from_id is blank" do
|
||||
db = open_db
|
||||
allow(message_harness).to receive(:normalize_node_id) do |_db, ref|
|
||||
if ref.nil? || ref.to_s.strip.empty?
|
||||
"!aabbccdd"
|
||||
else
|
||||
parts = message_harness.canonical_node_parts(ref)
|
||||
parts ? parts[0] : nil
|
||||
end
|
||||
end
|
||||
message_harness.insert_message(db, {
|
||||
"id" => 7001,
|
||||
"from_id" => "",
|
||||
"to_id" => "^all",
|
||||
"text" => "hello",
|
||||
"channel" => 0,
|
||||
})
|
||||
expect(db.get_first_value("SELECT from_id FROM messages WHERE id = 7001")).to eq("!aabbccdd")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
|
||||
it "rewrites a !-prefixed from_id when canonical resolution differs" do
|
||||
db = open_db
|
||||
allow(message_harness).to receive(:normalize_node_id) do |_db, ref|
|
||||
next nil if ref.nil? || ref.to_s.strip.empty?
|
||||
if ref.to_s.start_with?("!")
|
||||
"!11111111"
|
||||
else
|
||||
parts = message_harness.canonical_node_parts(ref)
|
||||
parts ? parts[0] : nil
|
||||
end
|
||||
end
|
||||
message_harness.insert_message(db, {
|
||||
"id" => 7002,
|
||||
"from_id" => "!aabbccdd",
|
||||
"to_id" => "^all",
|
||||
"text" => "hi",
|
||||
"channel" => 0,
|
||||
})
|
||||
expect(db.get_first_value("SELECT from_id FROM messages WHERE id = 7002")).to eq("!11111111")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
|
||||
it "uses canonical_to_id when the raw to_id is blank" do
|
||||
db = open_db
|
||||
allow(message_harness).to receive(:normalize_node_id) do |_db, ref|
|
||||
if ref.nil? || ref.to_s.strip.empty?
|
||||
"!ddddeeee"
|
||||
else
|
||||
parts = message_harness.canonical_node_parts(ref)
|
||||
parts ? parts[0] : nil
|
||||
end
|
||||
end
|
||||
message_harness.insert_message(db, {
|
||||
"id" => 7003,
|
||||
"from_id" => "!aabbccdd",
|
||||
"to_id" => "",
|
||||
"text" => "yo",
|
||||
"channel" => 0,
|
||||
})
|
||||
expect(db.get_first_value("SELECT to_id FROM messages WHERE id = 7003")).to eq("!ddddeeee")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
|
||||
it "rewrites a !-prefixed to_id when canonical resolution differs" do
|
||||
db = open_db
|
||||
allow(message_harness).to receive(:normalize_node_id) do |_db, ref|
|
||||
next nil if ref.nil? || ref.to_s.strip.empty?
|
||||
if ref.to_s == "!ffffeeee"
|
||||
"!22222222"
|
||||
else
|
||||
parts = message_harness.canonical_node_parts(ref)
|
||||
parts ? parts[0] : nil
|
||||
end
|
||||
end
|
||||
message_harness.insert_message(db, {
|
||||
"id" => 7004,
|
||||
"from_id" => "!aabbccdd",
|
||||
"to_id" => "!ffffeeee",
|
||||
"text" => "msg",
|
||||
"channel" => 0,
|
||||
})
|
||||
expect(db.get_first_value("SELECT to_id FROM messages WHERE id = 7004")).to eq("!22222222")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
|
||||
describe "#insert_message — reply_id and emoji updates on existing rows" do
|
||||
include_context "with isolated db"
|
||||
|
||||
let(:msg_update_harness) do
|
||||
Class.new do
|
||||
include PotatoMesh::App::DataProcessing
|
||||
include PotatoMesh::App::Helpers
|
||||
|
||||
def debug_log(*); end
|
||||
|
||||
def warn_log(*); end
|
||||
|
||||
def with_busy_retry
|
||||
yield
|
||||
end
|
||||
|
||||
def update_prometheus_metrics(*); end
|
||||
|
||||
def prom_report_ids
|
||||
[]
|
||||
end
|
||||
|
||||
def private_mode?
|
||||
false
|
||||
end
|
||||
|
||||
def resolve_protocol(*)
|
||||
"meshtastic"
|
||||
end
|
||||
|
||||
def normalize_node_id(_db, ref)
|
||||
parts = canonical_node_parts(ref)
|
||||
parts ? parts[0] : nil
|
||||
end
|
||||
|
||||
def ensure_unknown_node(*); end
|
||||
|
||||
def touch_node_last_seen(*); end
|
||||
end.new
|
||||
end
|
||||
|
||||
it "updates reply_id when the existing row references a different reply_id" do
|
||||
db = open_db
|
||||
base = {
|
||||
"from_id" => "!aabbccdd",
|
||||
"to_id" => "^all",
|
||||
"text" => "reply",
|
||||
"channel" => 0,
|
||||
"reply_id" => 100,
|
||||
}
|
||||
msg_update_harness.insert_message(db, base.merge("id" => 8001))
|
||||
msg_update_harness.insert_message(db, base.merge("id" => 8001, "reply_id" => 200))
|
||||
expect(db.get_first_value("SELECT reply_id FROM messages WHERE id = 8001")).to eq(200)
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
|
||||
it "fills in an emoji when the existing row had none" do
|
||||
db = open_db
|
||||
base = {
|
||||
"from_id" => "!aabbccdd",
|
||||
"to_id" => "^all",
|
||||
"text" => "thanks",
|
||||
"channel" => 0,
|
||||
}
|
||||
msg_update_harness.insert_message(db, base.merge("id" => 8002))
|
||||
msg_update_harness.insert_message(db, base.merge("id" => 8002, "emoji" => ":thumbsup:"))
|
||||
expect(db.get_first_value("SELECT emoji FROM messages WHERE id = 8002")).to eq(":thumbsup:")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
|
||||
describe "#insert_message — ConstraintException recovery" do
|
||||
include_context "with isolated db"
|
||||
|
||||
let(:fb_harness) do
|
||||
Class.new do
|
||||
include PotatoMesh::App::DataProcessing
|
||||
include PotatoMesh::App::Helpers
|
||||
|
||||
def debug_log(*); end
|
||||
|
||||
def warn_log(*); end
|
||||
|
||||
def with_busy_retry
|
||||
yield
|
||||
end
|
||||
|
||||
def update_prometheus_metrics(*); end
|
||||
|
||||
def prom_report_ids
|
||||
[]
|
||||
end
|
||||
|
||||
def private_mode?
|
||||
false
|
||||
end
|
||||
|
||||
def resolve_protocol(*)
|
||||
"meshtastic"
|
||||
end
|
||||
|
||||
def normalize_node_id(_db, ref)
|
||||
parts = canonical_node_parts(ref)
|
||||
parts ? parts[0] : nil
|
||||
end
|
||||
|
||||
def ensure_unknown_node(*); end
|
||||
|
||||
def touch_node_last_seen(*); end
|
||||
end.new
|
||||
end
|
||||
|
||||
let(:base_msg) do
|
||||
{
|
||||
"from_id" => "!aabbccdd",
|
||||
"to_id" => "^all",
|
||||
"channel" => 0,
|
||||
}
|
||||
end
|
||||
|
||||
# The fallback path is taken when the SELECT-before-INSERT misses but the
|
||||
# INSERT itself trips the PK constraint — i.e., a concurrent ingestor has
|
||||
# already inserted the row in between. We simulate that race by seeding a
|
||||
# row, then forcing the existing-row SELECT to return nil so the INSERT
|
||||
# path runs and trips the constraint deterministically.
|
||||
it "applies fallback updates when INSERT trips a constraint violation" do
|
||||
db = open_db
|
||||
fb_harness.insert_message(db, base_msg.merge("id" => 9001, "text" => "first", "ingestor" => "!a"))
|
||||
|
||||
allow(db).to receive(:get_first_row).and_wrap_original do |original, sql, *args|
|
||||
if sql.include?("SELECT from_id, to_id, text, encrypted, lora_freq")
|
||||
nil
|
||||
else
|
||||
original.call(sql, *args)
|
||||
end
|
||||
end
|
||||
|
||||
fb_harness.insert_message(db, base_msg.merge("id" => 9001, "text" => "second", "ingestor" => "!b"))
|
||||
|
||||
expect(db.get_first_value("SELECT text FROM messages WHERE id = 9001")).to eq("second")
|
||||
# First-write-wins for ingestor: the existing value (!a) is preserved.
|
||||
expect(db.get_first_value("SELECT ingestor FROM messages WHERE id = 9001")).to eq("!a")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
|
||||
it "applies decrypted-precedence overrides during fallback" do
|
||||
db = open_db
|
||||
fb_harness.insert_message(db, base_msg.merge(
|
||||
"id" => 9002,
|
||||
"encrypted" => "BLOB",
|
||||
"text" => nil,
|
||||
))
|
||||
|
||||
allow(db).to receive(:get_first_row).and_wrap_original do |original, sql, *args|
|
||||
if sql.include?("SELECT from_id, to_id, text, encrypted, lora_freq")
|
||||
nil
|
||||
else
|
||||
original.call(sql, *args)
|
||||
end
|
||||
end
|
||||
|
||||
fb_harness.insert_message(db, base_msg.merge(
|
||||
"id" => 9002,
|
||||
"text" => "decrypted",
|
||||
"lora_freq" => 869525,
|
||||
"modem_preset" => "MEDIUM_SLOW",
|
||||
"channel_name" => "LongFast",
|
||||
))
|
||||
|
||||
expect(db.get_first_value("SELECT text FROM messages WHERE id = 9002")).to eq("decrypted")
|
||||
expect(db.get_first_value("SELECT lora_freq FROM messages WHERE id = 9002")).to eq(869525)
|
||||
expect(db.get_first_value("SELECT modem_preset FROM messages WHERE id = 9002")).to eq("MEDIUM_SLOW")
|
||||
ensure
|
||||
db&.close
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
Reference in New Issue
Block a user