fix: prevent NPM tunnel collisions and false app health restarts
This commit is contained in:
@@ -0,0 +1,158 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Migrate the known NPM/LND tunnel collision, without touching wallet services.
|
||||
|
||||
Runs as the rootless app owner before orchestrator startup. Only the narrow
|
||||
legacy web-tunnel profile is accepted. Unknown custom routing fails closed.
|
||||
"""
|
||||
import ipaddress
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import re
|
||||
import socket
|
||||
import subprocess
|
||||
import tempfile
|
||||
|
||||
|
||||
def command(*args, input=None):
|
||||
result = subprocess.run(args, input=input, text=True, capture_output=True, timeout=45)
|
||||
if result.returncode:
|
||||
# Commands may read private files. Never print their captured output.
|
||||
raise RuntimeError(f'{args[0]} operation failed (exit {result.returncode})')
|
||||
return result.stdout
|
||||
|
||||
|
||||
def plan(drop, rules):
|
||||
"""Return a conservative migration, or None for absent/already fixed mapping."""
|
||||
matches = re.findall(r'^PublishPort=([0-9.]+):18080:80/tcp$', drop, re.M)
|
||||
if not matches:
|
||||
return None
|
||||
if len(matches) != 1:
|
||||
raise ValueError('ambiguous NPM tunnel mapping')
|
||||
destination = str(ipaddress.IPv4Address(matches[0]))
|
||||
peer_match = re.search(r'ip saddr ([0-9.]+) ip daddr ' + re.escape(destination)
|
||||
+ r' tcp dport \{ 18080, 18443 \} accept', rules)
|
||||
if not peer_match:
|
||||
raise ValueError('unrecognized NPM tunnel firewall; manual review required')
|
||||
peer = str(ipaddress.IPv4Address(peer_match[1]))
|
||||
# Match the entire old profile, not just a substring in an arbitrary firewall.
|
||||
old = f'''table inet web_tunnel {{
|
||||
chain input {{
|
||||
type filter hook input priority -10; policy accept;
|
||||
iifname != "wg-web" return
|
||||
ct state established,related accept
|
||||
ip saddr {peer} icmp type echo-request accept
|
||||
ip saddr {peer} ip daddr {destination} tcp dport {{ 18080, 18443 }} accept
|
||||
counter drop
|
||||
}}
|
||||
chain forward {{
|
||||
type filter hook forward priority -10; policy accept;
|
||||
iifname "wg-web" counter drop
|
||||
oifname "wg-web" counter drop
|
||||
}}
|
||||
}}'''
|
||||
if rules.split() != old.split():
|
||||
raise ValueError('custom NPM tunnel firewall differs; manual review required')
|
||||
if 'PublishPort='+destination+':18081:' in drop:
|
||||
raise ValueError('replacement port already configured')
|
||||
new_rules = rules.replace('table inet web_tunnel {', f'''table inet web_tunnel {{
|
||||
# Preserve incoming HTTP while keeping LND REST's port free.
|
||||
chain prerouting {{
|
||||
type nat hook prerouting priority dstnat; policy accept;
|
||||
iifname "wg-web" ip saddr {peer} ip daddr {destination} tcp dport 18080 redirect to :18081
|
||||
}}''', 1).replace('tcp dport { 18080, 18443 } accept',
|
||||
'tcp dport { 18081, 18443 } accept')
|
||||
return (drop.replace(f'PublishPort={destination}:18080:80/tcp',
|
||||
f'PublishPort={destination}:18081:80/tcp'), new_rules, destination)
|
||||
|
||||
|
||||
def atomic_user(path, content):
|
||||
with tempfile.NamedTemporaryFile(mode='w', dir=path.parent, delete=False) as f:
|
||||
tmp = Path(f.name)
|
||||
os.fchmod(f.fileno(), 0o600)
|
||||
f.write(content)
|
||||
f.flush()
|
||||
os.fsync(f.fileno())
|
||||
os.replace(tmp, path)
|
||||
|
||||
|
||||
def root_write(path, content):
|
||||
# Stage next to the destination; rename makes the config update atomic.
|
||||
staged = str(path)+'.archy-npm-migration'
|
||||
command('sudo', '-n', 'tee', staged, input=content)
|
||||
command('sudo', '-n', 'chmod', '600', staged)
|
||||
command('sudo', '-n', 'mv', '--', staged, str(path))
|
||||
|
||||
|
||||
def main():
|
||||
drop = Path.home()/'.config/containers/systemd/nginx-proxy-manager.container.d/web-tunnel.conf'
|
||||
rules_path = Path('/etc/wireguard/wg-web.nft')
|
||||
state = Path.home()/'.local/state/archipelago/npm-tunnel-migration'
|
||||
journal = state/'pending.json'
|
||||
recovered_active = None
|
||||
# Interrupted migrations are completed/rolled back before normal startup.
|
||||
if journal.exists():
|
||||
saved = json.loads(journal.read_text())
|
||||
command('systemctl', '--user', 'stop', 'nginx-proxy-manager.service')
|
||||
atomic_user(drop, saved['drop'])
|
||||
root_write(rules_path, saved['rules'])
|
||||
command('sudo', '-n', 'nft', '-f', '-', input='delete table inet web_tunnel\n'+saved['rules'])
|
||||
command('systemctl', '--user', 'daemon-reload')
|
||||
# Do not restart the colliding old configuration before reapplying.
|
||||
recovered_active = saved.get('was_active')
|
||||
journal.unlink()
|
||||
if not drop.exists():
|
||||
return
|
||||
old_drop = drop.read_text()
|
||||
if not re.search(r'^PublishPort=[0-9.]+:18080:80/tcp$', old_drop, re.M):
|
||||
return
|
||||
old_rules = command('sudo', '-n', 'cat', str(rules_path))
|
||||
new_drop, new_rules, destination = plan(old_drop, old_rules)
|
||||
# A free, assigned replacement is required; do not guess another port.
|
||||
with socket.socket() as probe:
|
||||
probe.bind((destination, 18081))
|
||||
# The route must be persistent and loaded by the tunnel's startup contract.
|
||||
wg = command('sudo', '-n', 'grep', '-E', r'^(PreUp|PostDown)\s*=', '/etc/wireguard/wg-web.conf')
|
||||
if 'PreUp = nft -f /etc/wireguard/wg-web.nft' not in wg or 'PostDown = nft delete table inet web_tunnel' not in wg:
|
||||
raise ValueError('unrecognized tunnel lifecycle; manual review required')
|
||||
active = command('sudo', '-n', 'nft', 'list', 'table', 'inet', 'web_tunnel')
|
||||
# Reject live-only rule changes instead of silently discarding them. nft
|
||||
# canonicalizes priority names and adds counter values when listing rules.
|
||||
def normalized(text):
|
||||
text = re.sub(r'counter packets \d+ bytes \d+', 'counter', text)
|
||||
return text.replace('priority filter - 10', 'priority -10').split()
|
||||
if normalized(active) != normalized(old_rules):
|
||||
raise ValueError('live tunnel rules differ from persistent config; review required')
|
||||
transaction = 'delete table inet web_tunnel\n'+new_rules
|
||||
command('sudo', '-n', 'nft', '--check', '-f', '-', input=transaction)
|
||||
state.mkdir(parents=True, exist_ok=True, mode=0o700)
|
||||
os.chmod(state, 0o700)
|
||||
was_active = recovered_active or command('systemctl', '--user', 'show', 'nginx-proxy-manager.service', '--property=ActiveState', '--value').strip()
|
||||
saved = json.dumps({'drop': old_drop, 'rules': old_rules, 'was_active': was_active})
|
||||
atomic_user(state/'before.json', saved)
|
||||
atomic_user(journal, saved)
|
||||
try:
|
||||
command('systemctl', '--user', 'stop', 'nginx-proxy-manager.service')
|
||||
atomic_user(drop, new_drop)
|
||||
root_write(rules_path, new_rules)
|
||||
command('sudo', '-n', 'nft', '-f', '-', input=transaction)
|
||||
command('systemctl', '--user', 'daemon-reload')
|
||||
if was_active in ('active', 'activating', 'reloading', 'failed'):
|
||||
command('systemctl', '--user', 'restart', 'nginx-proxy-manager.service')
|
||||
journal.unlink()
|
||||
except Exception:
|
||||
atomic_user(drop, old_drop)
|
||||
root_write(rules_path, old_rules)
|
||||
command('sudo', '-n', 'nft', '-f', '-', input='delete table inet web_tunnel\n'+old_rules)
|
||||
command('systemctl', '--user', 'daemon-reload')
|
||||
# Keep the journal if rollback fails so the next startup retries it.
|
||||
journal.unlink()
|
||||
raise
|
||||
print('NPM tunnel port repaired; original configuration backed up; native services unchanged')
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
try:
|
||||
main()
|
||||
except Exception as error:
|
||||
raise SystemExit('NPM tunnel migration requires attention: '+str(error)) from None
|
||||
@@ -0,0 +1,137 @@
|
||||
import importlib.util
|
||||
from pathlib import Path
|
||||
import unittest
|
||||
from unittest.mock import patch, MagicMock
|
||||
import tempfile
|
||||
import json
|
||||
|
||||
spec=importlib.util.spec_from_file_location('repair', Path(__file__).parents[1]/'repair-npm-tunnel.py')
|
||||
m=importlib.util.module_from_spec(spec)
|
||||
spec.loader.exec_module(m)
|
||||
DROP='[Container]\nPublishPort=10.77.0.2:18080:80/tcp\nPublishPort=10.77.0.2:18443:443/tcp\n'
|
||||
RULES='''table inet web_tunnel {
|
||||
chain input {
|
||||
type filter hook input priority -10; policy accept;
|
||||
iifname != "wg-web" return
|
||||
ct state established,related accept
|
||||
ip saddr 10.77.0.1 icmp type echo-request accept
|
||||
ip saddr 10.77.0.1 ip daddr 10.77.0.2 tcp dport { 18080, 18443 } accept
|
||||
counter drop
|
||||
}
|
||||
chain forward {
|
||||
type filter hook forward priority -10; policy accept;
|
||||
iifname "wg-web" counter drop
|
||||
oifname "wg-web" counter drop
|
||||
}
|
||||
}'''
|
||||
|
||||
class Plan(unittest.TestCase):
|
||||
def test_preserves_peer_https_and_restricts_redirect(self):
|
||||
drop,rules,dst=m.plan(DROP,RULES)
|
||||
self.assertEqual(dst,'10.77.0.2')
|
||||
self.assertIn('PublishPort=10.77.0.2:18081:80/tcp',drop)
|
||||
self.assertIn('PublishPort=10.77.0.2:18443:443/tcp',drop)
|
||||
self.assertIn('iifname "wg-web" ip saddr 10.77.0.1 ip daddr 10.77.0.2 tcp dport 18080 redirect to :18081',rules)
|
||||
self.assertNotIn('tcp dport { 18080, 18443 } accept',rules)
|
||||
self.assertIn('iifname "wg-web" counter drop',rules)
|
||||
def test_idempotent(self):
|
||||
d,r,_=m.plan(DROP,RULES)
|
||||
self.assertIsNone(m.plan(d,r))
|
||||
def test_standard_fresh_install_untouched(self):
|
||||
self.assertIsNone(m.plan('[Container]\nPublishPort=127.0.0.1:8081:81/tcp',''))
|
||||
def test_no_hardcoded_deployment_address(self):
|
||||
d,r,dst=m.plan(DROP.replace('10.77.0.','10.55.0.'),RULES.replace('10.77.0.','10.55.0.'))
|
||||
self.assertEqual(dst,'10.55.0.2');self.assertIn('ip saddr 10.55.0.1',r)
|
||||
def test_custom_firewall_preserved(self):
|
||||
for r in [RULES+'\ntable inet extra {}',RULES.replace('counter drop','accept'),RULES.replace('wg-web','wg-custom')]:
|
||||
with self.assertRaises(ValueError): m.plan(DROP,r)
|
||||
def test_ambiguous_mapping(self):
|
||||
with self.assertRaises(ValueError): m.plan(DROP+DROP,RULES)
|
||||
def test_wrong_destination(self):
|
||||
with self.assertRaises(ValueError): m.plan(DROP.replace('10.77.0.2','10.77.0.3'),RULES)
|
||||
def test_already_used_mapping(self):
|
||||
with self.assertRaises(ValueError): m.plan(DROP+'PublishPort=10.77.0.2:18081:80/tcp\n',RULES)
|
||||
|
||||
class Migration(unittest.TestCase):
|
||||
def setUp(self):
|
||||
self.temp=tempfile.TemporaryDirectory()
|
||||
self.addCleanup(self.temp.cleanup)
|
||||
self.home=Path(self.temp.name)
|
||||
self.drop=self.home/'.config/containers/systemd/nginx-proxy-manager.container.d/web-tunnel.conf'
|
||||
self.drop.parent.mkdir(parents=True)
|
||||
self.drop.write_text(DROP)
|
||||
self.calls=[];self.rules=RULES;self.fail=None
|
||||
self.state=self.home/'.local/state/archipelago/npm-tunnel-migration'
|
||||
def command(self,*args,input=None):
|
||||
self.calls.append((args,input))
|
||||
if self.fail and self.fail(args):
|
||||
self.fail=None
|
||||
raise RuntimeError('injected failure')
|
||||
if 'cat' in args or ('list' in args and 'nft' in args): return self.rules
|
||||
if 'grep' in args: return 'PreUp = nft -f /etc/wireguard/wg-web.nft\nPostDown = nft delete table inet web_tunnel'
|
||||
if 'show' in args: return 'active'
|
||||
return ''
|
||||
def run_migration(self):
|
||||
with patch.object(Path,'home',return_value=self.home),patch.object(m,'command',side_effect=self.command),patch.object(m,'root_write') as write,patch.object(m.socket,'socket'):
|
||||
m.main()
|
||||
return write
|
||||
def test_success_and_second_run_noop(self):
|
||||
write=self.run_migration()
|
||||
self.assertIn(':18081:80/tcp',self.drop.read_text())
|
||||
self.assertEqual(json.loads((self.state/'before.json').read_text())['drop'],DROP)
|
||||
self.assertFalse((self.state/'pending.json').exists())
|
||||
self.assertEqual(write.call_count,1)
|
||||
self.calls.clear();self.run_migration();self.assertEqual(self.calls,[])
|
||||
def test_no_native_service_commands(self):
|
||||
self.run_migration()
|
||||
for args,_ in self.calls:
|
||||
self.assertNotIn('lnd.service',args);self.assertNotIn('bitcoin-core.service',args)
|
||||
def test_validation_failure_does_not_stop_or_write(self):
|
||||
self.fail=lambda a:'--check' in a
|
||||
with self.assertRaises(RuntimeError):self.run_migration()
|
||||
self.assertEqual(self.drop.read_text(),DROP)
|
||||
self.assertFalse(self.state.exists())
|
||||
self.assertFalse(any('stop' in a for a,_ in self.calls))
|
||||
def test_apply_failure_restores_files_and_firewall(self):
|
||||
self.fail=lambda a:'nft' in a and '-f' in a and '--check' not in a
|
||||
with self.assertRaises(RuntimeError):self.run_migration()
|
||||
self.assertEqual(self.drop.read_text(),DROP)
|
||||
self.assertFalse((self.state/'pending.json').exists())
|
||||
self.assertTrue(any(v=='delete table inet web_tunnel\n'+RULES for _,v in self.calls))
|
||||
def test_crash_journal_recovers_and_retries(self):
|
||||
self.state.mkdir(parents=True)
|
||||
(self.state/'pending.json').write_text(json.dumps({'drop':DROP,'rules':RULES,'was_active':'active'}))
|
||||
self.drop.write_text(m.plan(DROP,RULES)[0])
|
||||
self.run_migration()
|
||||
self.assertIn(':18081:80/tcp',self.drop.read_text())
|
||||
self.assertFalse((self.state/'pending.json').exists())
|
||||
def test_fresh_install_executes_no_commands(self):
|
||||
self.drop.unlink();self.run_migration();self.assertEqual(self.calls,[])
|
||||
def test_busy_replacement_port_does_not_mutate(self):
|
||||
with patch.object(Path,'home',return_value=self.home),patch.object(m,'command',side_effect=self.command),patch.object(m.socket,'socket') as socket:
|
||||
socket.return_value.__enter__.return_value.bind.side_effect=OSError('in use')
|
||||
with self.assertRaises(OSError):m.main()
|
||||
self.assertFalse(self.state.exists())
|
||||
self.assertFalse(any('stop' in a for a,_ in self.calls))
|
||||
def test_failed_rollback_keeps_recovery_journal(self):
|
||||
with patch.object(Path,'home',return_value=self.home),patch.object(m,'command',side_effect=self.command),patch.object(m,'root_write',side_effect=RuntimeError('write failed')),patch.object(m.socket,'socket'):
|
||||
with self.assertRaises(RuntimeError):m.main()
|
||||
self.assertTrue((self.state/'pending.json').exists())
|
||||
self.assertEqual(self.drop.read_text(),DROP)
|
||||
def test_stopped_app_is_not_started(self):
|
||||
original=self.command
|
||||
def stopped(*args,input=None):
|
||||
return 'inactive' if 'show' in args else original(*args,input=input)
|
||||
with patch.object(Path,'home',return_value=self.home),patch.object(m,'command',side_effect=stopped),patch.object(m,'root_write'),patch.object(m.socket,'socket'):
|
||||
m.main()
|
||||
self.assertFalse(any('restart' in a for a,_ in self.calls))
|
||||
def test_live_only_firewall_changes_are_not_discarded(self):
|
||||
original=self.command
|
||||
def different(*args,input=None):
|
||||
result=original(*args,input=input)
|
||||
return result+' table inet custom {}' if 'list' in args else result
|
||||
with patch.object(Path,'home',return_value=self.home),patch.object(m,'command',side_effect=different),patch.object(m.socket,'socket'):
|
||||
with self.assertRaises(ValueError):m.main()
|
||||
self.assertEqual(self.drop.read_text(),DROP)
|
||||
|
||||
if __name__=='__main__': unittest.main()
|
||||
Reference in New Issue
Block a user