Files
l5y a676a2997e data,web: Meshtastic waypoints as a first-class POI layer (#874)
* data,web: Meshtastic waypoints as a first-class POI layer

* web: re-roll waypoint UI to signed-off variants + node-page section

* web: reconcile stats-pushdown plan spec with the five-table W9 umbrella

* web: legend waypoint sample follows the 1c-B teardrop marker
2026-07-29 18:22:08 +02:00

424 lines
17 KiB
Ruby

# 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
require "spec_helper"
RSpec.describe PotatoMesh::App::PubSub do
# Drop any globally registered subscribers so examples never leak state.
# Reference the module explicitly: nested groups rebind +described_class+.
after { PotatoMesh::App::PubSub.reset! }
describe "COLLECTIONS" do
it "lists exactly the seven dashboard ingest collections (waypoints joined, SPEC W8)" do
expect(described_class::COLLECTIONS).to eq(
%w[nodes messages positions telemetry neighbors traces waypoints],
)
end
it "suppresses exactly the message-grade collections in private mode (PS6/W3)" do
expect(described_class::PRIVATE_SUPPRESSED_COLLECTIONS).to eq(%w[messages waypoints])
end
end
describe ".subscribe" do
it "returns a Subscriber and tracks it" do
subscriber = described_class.subscribe
expect(subscriber).to be_a(described_class::Subscriber)
expect(described_class.subscriber_count).to eq(1)
end
it "raises CapacityError once MAX_SUBSCRIBERS is reached" do
described_class::MAX_SUBSCRIBERS.times { described_class.subscribe }
expect(described_class.subscriber_count).to eq(described_class::MAX_SUBSCRIBERS)
expect { described_class.subscribe }.to raise_error(described_class::CapacityError)
end
end
describe ".effective_max_subscribers" do
it "equals MAX_SUBSCRIBERS at the default pool size (96 - 32 = 64)" do
within_env("MAX_THREADS" => nil, "SSE_THREAD_RESERVE" => nil) do
expect(described_class.effective_max_subscribers).to eq(described_class::MAX_SUBSCRIBERS)
end
end
it "shrinks to the thread budget when the pool is smaller than the cap" do
within_env("MAX_THREADS" => "20", "SSE_THREAD_RESERVE" => "4") do
expect(described_class.effective_max_subscribers).to eq(16)
end
end
it "never goes negative when the pool is at or below the reserve" do
within_env("MAX_THREADS" => "4", "SSE_THREAD_RESERVE" => "8") do
expect(described_class.effective_max_subscribers).to eq(0)
end
end
it "caps subscribe at the clamped limit, then raises naming that limit" do
within_env("MAX_THREADS" => "10", "SSE_THREAD_RESERVE" => "8") do
expect(described_class.effective_max_subscribers).to eq(2)
2.times { described_class.subscribe }
expect { described_class.subscribe }.to raise_error(described_class::CapacityError, /subscriber limit \(2\)/)
end
end
end
describe ".unsubscribe" do
it "removes and closes the subscriber" do
subscriber = described_class.subscribe
described_class.unsubscribe(subscriber)
expect(described_class.subscriber_count).to eq(0)
expect(subscriber.closed?).to be(true)
end
it "is idempotent and nil-safe" do
subscriber = described_class.subscribe
described_class.unsubscribe(subscriber)
expect { described_class.unsubscribe(subscriber) }.not_to raise_error
expect { described_class.unsubscribe(nil) }.not_to raise_error
end
end
describe ".publish" do
it "delivers a change to every subscriber and returns the delivery count" do
a = described_class.subscribe
b = described_class.subscribe
expect(described_class.publish("nodes")).to eq(2)
expect(a.pending_count).to eq(1)
expect(b.pending_count).to eq(1)
end
it "ignores unknown collections without delivering" do
subscriber = described_class.subscribe
expect(described_class.publish("bogus")).to eq(0)
expect(subscriber.pending_count).to eq(0)
end
it "suppresses messages events in private mode (PS6)" do
subscriber = described_class.subscribe
expect(described_class.publish("messages", private_mode: true)).to eq(0)
expect(subscriber.pending_count).to eq(0)
end
it "suppresses waypoints events in private mode (SPEC W3 message-grade privacy)" do
subscriber = described_class.subscribe
expect(described_class.publish("waypoints", private_mode: true)).to eq(0)
expect(subscriber.pending_count).to eq(0)
end
it "delivers waypoints events when not private (seventh collection, SPEC W8)" do
subscriber = described_class.subscribe
expect(described_class.publish("waypoints")).to eq(1)
expect(subscriber.drain(timeout: 0.1)).to eq([{ collection: "waypoints", hint: nil }])
end
it "still delivers non-message collections in private mode" do
subscriber = described_class.subscribe
expect(described_class.publish("nodes", private_mode: true)).to eq(1)
expect(subscriber.pending_count).to eq(1)
end
it "delivers messages events when not private" do
subscriber = described_class.subscribe
expect(described_class.publish("messages")).to eq(1)
expect(subscriber.pending_count).to eq(1)
end
it "returns 0 when there are no subscribers" do
expect(described_class.publish("nodes")).to eq(0)
end
it "forwards the newest rx_time hint to subscribers" do
subscriber = described_class.subscribe
described_class.publish("messages", hint: 1_700_000_000)
expect(subscriber.drain(timeout: 0).first).to eq(
collection: "messages", hint: 1_700_000_000,
)
end
end
describe ".reset!" do
it "closes and clears every subscriber" do
a = described_class.subscribe
b = described_class.subscribe
described_class.reset!
expect(described_class.subscriber_count).to eq(0)
expect(a.closed?).to be(true)
expect(b.closed?).to be(true)
end
end
describe PotatoMesh::App::PubSub::Subscriber do
subject(:subscriber) { described_class.new }
describe "#deliver / #drain coalescing" do
it "collapses repeated writes to one collection into a single event" do
5.times { subscriber.deliver("messages", nil) }
expect(subscriber.pending_count).to eq(1)
events = subscriber.drain(timeout: 0)
expect(events).to eq([{ collection: "messages", hint: nil }])
end
it "keeps the newest (largest) hint across coalesced writes" do
subscriber.deliver("nodes", 5)
subscriber.deliver("nodes", 3)
subscriber.deliver("nodes", 9)
expect(subscriber.drain(timeout: 0)).to eq([{ collection: "nodes", hint: 9 }])
end
it "fills a missing hint from a later write" do
subscriber.deliver("nodes", nil)
subscriber.deliver("nodes", 7)
expect(subscriber.drain(timeout: 0)).to eq([{ collection: "nodes", hint: 7 }])
end
it "returns pending collections sorted by name, then clears them" do
subscriber.deliver("traces", nil)
subscriber.deliver("nodes", nil)
subscriber.deliver("messages", nil)
expect(subscriber.drain(timeout: 0).map { |e| e[:collection] }).to eq(
%w[messages nodes traces],
)
expect(subscriber.drain(timeout: 0.01)).to eq([])
end
end
describe "#drain blocking" do
it "returns an empty array after the timeout when idle" do
expect(subscriber.drain(timeout: 0.02)).to eq([])
end
it "wakes as soon as a change is delivered from another thread" do
producer = Thread.new do
sleep 0.02
subscriber.deliver("positions", 42)
end
events = subscriber.drain(timeout: 5)
producer.join
expect(events).to eq([{ collection: "positions", hint: 42 }])
end
end
describe "#drain settle window (LV6 cooldown)" do
it "coalesces a burst within the settle window into one drain" do
subscriber.deliver("positions", 1)
# The injected sleeper stands in for the cooldown window; more ingestors
# deliver the same/another collection while it "sleeps".
sleeper = lambda do |_seconds|
subscriber.deliver("positions", 2)
subscriber.deliver("telemetry", 5)
end
events = subscriber.drain(timeout: 0, settle: 1, sleeper: sleeper)
expect(events).to eq(
[
{ collection: "positions", hint: 2 },
{ collection: "telemetry", hint: 5 },
],
)
end
it "sleeps for the settle window only when a change is pending" do
slept = []
sleeper = ->(seconds) { slept << seconds }
# Idle heartbeat tick: nothing pending, so no settle sleep.
expect(subscriber.drain(timeout: 0, settle: 1, sleeper: sleeper)).to eq([])
expect(slept).to eq([])
# A pending change triggers exactly one settle sleep.
subscriber.deliver("nodes", nil)
expect(subscriber.drain(timeout: 0, settle: 1, sleeper: sleeper)).to eq(
[{ collection: "nodes", hint: nil }],
)
expect(slept).to eq([1])
end
it "does not sleep when settle is zero (cooldown disabled)" do
slept = []
sleeper = ->(seconds) { slept << seconds }
subscriber.deliver("nodes", nil)
subscriber.drain(timeout: 0, settle: 0, sleeper: sleeper)
expect(slept).to eq([])
end
end
describe "#close" do
it "ignores deliveries and drains immediately once closed" do
subscriber.close
expect(subscriber.closed?).to be(true)
subscriber.deliver("nodes", 1)
expect(subscriber.pending_count).to eq(0)
expect(subscriber.drain(timeout: 5)).to eq([])
end
it "flushes already-pending events after close" do
subscriber.deliver("nodes", 1)
subscriber.close
expect(subscriber.drain(timeout: 0)).to eq([{ collection: "nodes", hint: 1 }])
end
end
end
describe "publish-on-change from ingest routes (integration)" do
include Rack::Test::Methods
def app
PotatoMesh::Application
end
let(:auth) do
{ "CONTENT_TYPE" => "application/json", "HTTP_AUTHORIZATION" => "Bearer spec-token" }
end
around do |example|
original_token = ENV["API_TOKEN"]
original_private = ENV["PRIVATE"]
ENV["API_TOKEN"] = "spec-token"
ENV.delete("PRIVATE")
example.run
original_token.nil? ? ENV.delete("API_TOKEN") : ENV["API_TOKEN"] = original_token
original_private.nil? ? ENV.delete("PRIVATE") : ENV["PRIVATE"] = original_private
PotatoMesh::App::PubSub.reset!
end
# collection => [route, minimal valid body that writes no rows]
def ingest_routes
{
"nodes" => ["/api/nodes", "{}"],
"messages" => ["/api/messages", "[]"],
"positions" => ["/api/positions", "[]"],
"telemetry" => ["/api/telemetry", "[]"],
"neighbors" => ["/api/neighbors", "[]"],
"traces" => ["/api/traces", "[]"],
"waypoints" => ["/api/waypoints", "[]"],
}
end
it "publishes a thin per-collection event (PS3)" do
subscriber = PotatoMesh::App::PubSub.subscribe
# Use neighbors: it publishes exactly one collection. (The messages,
# positions, and telemetry routes additionally publish nodes because they
# advance a node's last_heard - covered by the dedicated examples below.)
post "/api/neighbors", "[]", auth
expect(last_response.status).to eq(201)
expect(subscriber.drain(timeout: 0.1)).to eq([{ collection: "neighbors", hint: nil }])
end
it "publishes nodes on a message ingest (#822 touches the author node)" do
allow(PotatoMesh::App::PubSub).to receive(:publish).and_call_original
post "/api/messages", "[]", auth
expect(last_response.status).to eq(201)
expect(PotatoMesh::App::PubSub).to have_received(:publish).with("messages", private_mode: false)
expect(PotatoMesh::App::PubSub).to have_received(:publish).with("nodes", private_mode: false)
end
# A positions / telemetry ingest advances the affected node's last_heard
# server-side (touch_node_last_seen), but unlike the messages route it
# historically published only its own collection, so the live dashboard
# fetched only that collection and never re-pulled the node row, leaving
# the node table "last seen" stale until the safety poll. These routes now
# also publish "nodes" (mirroring #822) so the changed node last_heard
# refreshes and flashes live.
it "publishes nodes on a positions ingest (advances the node last_heard)" do
allow(PotatoMesh::App::PubSub).to receive(:publish).and_call_original
post "/api/positions", "[]", auth
expect(last_response.status).to eq(201)
expect(PotatoMesh::App::PubSub).to have_received(:publish).with("positions", private_mode: false)
expect(PotatoMesh::App::PubSub).to have_received(:publish).with("nodes", private_mode: false)
end
it "publishes nodes on a telemetry ingest (advances the node last_heard)" do
allow(PotatoMesh::App::PubSub).to receive(:publish).and_call_original
post "/api/telemetry", "[]", auth
expect(last_response.status).to eq(201)
expect(PotatoMesh::App::PubSub).to have_received(:publish).with("telemetry", private_mode: false)
expect(PotatoMesh::App::PubSub).to have_received(:publish).with("nodes", private_mode: false)
end
# Neighbors / traces also advance last_heard, but SPEC VF3 keeps them out
# of the flash scope, so they deliberately do NOT publish "nodes" (their
# last_heard refresh is surfaced silently by the safety poll, never flashed).
it "does not publish nodes on a neighbors or traces ingest (VF3 boundary)" do
allow(PotatoMesh::App::PubSub).to receive(:publish).and_call_original
post "/api/neighbors", "[]", auth
expect(last_response.status).to eq(201)
post "/api/traces", "[]", auth
expect(last_response.status).to eq(201)
expect(PotatoMesh::App::PubSub).not_to have_received(:publish).with("nodes", private_mode: false)
end
# Waypoints are on the FLASHING side of the live-update boundary (SPEC W8
# as re-rolled): a waypoint ingest advances the author's last_heard, so the
# route also publishes "nodes" — mirroring positions/messages — and the
# dashboard fades the pin and flashes the author with a fresh "last seen".
it "publishes nodes on a waypoints ingest (author flash, W8 re-roll)" do
allow(PotatoMesh::App::PubSub).to receive(:publish).and_call_original
post "/api/waypoints", "[]", auth
expect(last_response.status).to eq(201)
expect(PotatoMesh::App::PubSub).to have_received(:publish).with("waypoints", private_mode: false)
expect(PotatoMesh::App::PubSub).to have_received(:publish).with("nodes", private_mode: false)
end
it "publishes on every ingest route" do
allow(PotatoMesh::App::PubSub).to receive(:publish).and_call_original
ingest_routes.each_value do |(route, body)|
post route, body, auth
expect(last_response.status).to eq(201)
end
ingest_routes.each_key do |collection|
# `nodes` is published by both the nodes route and the messages route
# (#822), so assert at-least-once rather than exactly-once.
expect(PotatoMesh::App::PubSub).to have_received(:publish)
.with(collection, private_mode: false).at_least(:once)
end
end
it "coalesces bursts of one collection into a single pending event" do
subscriber = PotatoMesh::App::PubSub.subscribe
# neighbors publishes exactly one collection (positions/telemetry now
# also publish nodes for the last_heard hook-in), so it isolates the
# per-collection coalescing under test.
5.times do
post "/api/neighbors", "[]", auth
expect(last_response.status).to eq(201)
end
expect(subscriber.pending_count).to eq(1)
expect(subscriber.drain(timeout: 0.1)).to eq([{ collection: "neighbors", hint: nil }])
end
end
# Temporarily set/clear ENV keys for the duration of the block, restoring the
# prior values (including "previously unset") afterwards.
def within_env(values)
original = values.transform_values { |_| :__unset__ }
values.each do |key, value|
original[key] = ENV.key?(key) ? ENV[key] : :__unset__
value.nil? ? ENV.delete(key) : ENV[key] = value
end
yield
ensure
original.each do |key, value|
value == :__unset__ ? ENV.delete(key) : ENV[key] = value
end
end
end