Files
spothole/webserver/handlers/api/addspot.py
T

322 lines
16 KiB
Python

import asyncio
import hmac
import logging
import re
import threading
from typing import Any
import requests
import tornado
from tornado import httputil
from tornado.ioloop import IOLoop
from tornado.web import Application
from core.activity_utils import get_ref_regex_for_activity
from core.config import (
ALLOW_SPOTTING,
ALLOW_UPSTREAM_SPOTTING,
API_KEYS,
PROTECT_SPOT_SUBMISSION,
RECAPTCHA_SECRET_KEY,
)
from core.constants import UNKNOWN_BAND
from core.utils import infer_band_from_freq, safe_json_dumps
from data.spot import Spot
from providers.spot.dxcluster import DXCluster
from providers.spot.spot_provider import SpotProvider, SpotSubmissionError
from providers.spot.tiles import Tiles
logger = logging.getLogger(__name__)
RECAPTCHA_VERIFY_URL = "https://www.google.com/recaptcha/api/siteverify"
class APISpotHandler(tornado.web.RequestHandler):
"""API request handler for /api/v3/spot (POST)"""
def __init__(
self,
application: "Application",
request: httputil.HTTPServerRequest,
**kwargs: Any,
):
self._spots = None
self._spot_providers = None
super().__init__(application, request, **kwargs)
def initialize(self, spots, spot_providers=None):
self._spots = spots
self._spot_providers = spot_providers or []
async def post(self):
"""Handle the post request with spot data. This is async because it could go on to submit spots to providers
which could be slow, so we need to stop it blocking the whole web server"""
try:
# Reject if not allowed
if not ALLOW_SPOTTING:
self.set_status(401)
self.write(safe_json_dumps("Error - this server does not allow new spots to be added via the API."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Reject if format not json
if not self.request.headers.get("Content-Type", "").startswith("application/json"):
self.set_status(415)
self.write(safe_json_dumps("Error - request Content-Type must be application/json"))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Reject if request body is empty
post_data = self.request.body
if not post_data:
self.set_status(422)
self.write(safe_json_dumps("Error - request body is empty"))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Read in the request body as JSON
json_body = tornado.escape.json_decode(post_data)
# Extract the "spot" and "handling" sub-objects from the request body
spot_data = json_body.get("spot", {})
handling = json_body.get("handling", {})
# Extract individual parameters that say how this spot should be handled by the server
upstream_provider_names = handling.get("upstream_providers", None) or []
upstream_credentials = handling.get("upstream_credentials", None) or {}
submit_upstream = len(upstream_provider_names) > 0
captcha_token = handling.get("captcha_token", None)
# If spot submission is protected, the client must either provide a valid API key in the request header, or
# a valid CAPTCHA token in the request body. API keys are how we allow trusted third-party clients to submit
# spots, as they can't solve a CAPTCHA.
if PROTECT_SPOT_SUBMISSION:
api_key = self.request.headers.get("X-API-Key", "")
if api_key:
if not self._is_valid_api_key(api_key):
self.set_status(401)
self.write(safe_json_dumps("Error - API key not recognised."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
elif captcha_token and RECAPTCHA_SECRET_KEY:
if not await IOLoop.current().run_in_executor(None, self._verify_recaptcha, captcha_token):
self.set_status(422)
self.write(safe_json_dumps("Error - CAPTCHA verification failed."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
else:
self.set_status(401)
if RECAPTCHA_SECRET_KEY:
message = (
"Error - this server requires either an API key or a CAPTCHA token for spot submission."
)
else:
message = "Error - this server requires an API key for spot submission."
self.write(safe_json_dumps(message))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Convert spot field to a Spot object
spot = Spot(**spot_data)
# Reject if no timestamp, frequency, dx_call or de_call
if not spot.time or not spot.dx_call or not spot.freq or not spot.de_call:
self.set_status(422)
self.write(
safe_json_dumps("Error - 'time', 'dx_call', 'freq' and 'de_call' must be provided as a minimum.")
)
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Reject invalid-looking callsigns
if not re.match(r"^[A-Za-z0-9/\-]*$", spot.dx_call):
self.set_status(422)
self.write(safe_json_dumps(f"Error - '{spot.dx_call}' does not look like a valid callsign."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
if not re.match(r"^[A-Za-z0-9/\-]*$", spot.de_call):
self.set_status(422)
self.write(safe_json_dumps(f"Error - '{spot.de_call}' does not look like a valid callsign."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Reject if frequency not in a known band
if infer_band_from_freq(spot.freq) == UNKNOWN_BAND:
self.set_status(422)
self.write(safe_json_dumps(f"Error - Frequency of {spot.freq / 1000.0!s}kHz is not in a known band."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Reject if grid formatting incorrect
if spot.dx_grid and not re.match(
r"^([A-R]{2}[0-9]{2}[A-X]{2}[0-9]{2}[A-X]{2}|[A-R]{2}[0-9]{2}[A-X]{2}[0-9]{2}|[A-R]{2}[0-9]{2}[A-X]{2}|[A-R]{2}[0-9]{2})$",
spot.dx_grid.upper(),
):
self.set_status(422)
self.write(safe_json_dumps(f"Error - '{spot.dx_grid}' does not look like a valid Maidenhead grid."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Reject if any activity ref format is incorrect for its activity
for activity_ref in spot.activity_refs:
ref_regex = get_ref_regex_for_activity(activity_ref.activity) if activity_ref.activity else None
if activity_ref.id and ref_regex and not re.match(ref_regex, activity_ref.id):
self.set_status(422)
self.write(
safe_json_dumps(
f"Error - '{activity_ref.id}' does not look like a valid reference for {activity_ref.activity}."
)
)
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Reject upstream submission if not permitted
if submit_upstream and not ALLOW_UPSTREAM_SPOTTING:
self.set_status(403)
self.write(safe_json_dumps("Error - this server does not allow upstream spot submission."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Validate upstream submission requirements for each requested provider
for upstream_provider_name in upstream_provider_names:
provider = self._find_provider(upstream_provider_name, spot.activities)
is_cluster = isinstance(provider, DXCluster)
is_tiles = isinstance(provider, Tiles)
if not spot.activity_refs and not is_tiles and not is_cluster:
self.set_status(422)
self.write(
safe_json_dumps(
f"Error - an activity reference is required to submit upstream to {upstream_provider_name}."
)
)
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
if not spot.dx_grid and is_tiles:
self.set_status(422)
self.write(
safe_json_dumps("Error - a grid reference is required to submit upstream to Tiles on the Air.")
)
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
if not spot.mode and is_tiles:
self.set_status(422)
self.write(safe_json_dumps("Error - a mode is required to submit upstream to Tiles on the Air."))
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
return
# Submit upstream to all requested providers in parallel, collecting any warnings
results = await asyncio.gather(
*(
self._submit_upstream(name, spot, upstream_credentials.get(name, {}))
for name in upstream_provider_names
)
)
upstream_warnings = [w for w in results if w]
any_upstream_succeeded = len(upstream_warnings) < len(upstream_provider_names)
# If we successfully submitted the spot to at least one upstream provider, don't add it direct to Spothole,
# otherwise it will be a duplicate with what immediately comes back from the API. But if we weren't asked to
# send it upstream, or we were but every submission failed, we should still add it to our database anyway.
if not any_upstream_succeeded:
spot.source = "API"
spot.infer_missing()
self._spots.set(spot.id, spot)
logger.info(f"Spot of {spot.dx_call} by {spot.de_call} added to Spothole.")
if upstream_warnings:
saved_locally_note = "" if any_upstream_succeeded else " The spot was saved to Spothole only."
self.write(safe_json_dumps(f"Warning - {' '.join(upstream_warnings)}{saved_locally_note}"))
self.set_status(201)
else:
self.write(safe_json_dumps("OK"))
self.set_status(201)
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
except Exception:
logger.exception("Exception when handling client request to add spot API")
self.write(safe_json_dumps("Error - an internal server error occurred."))
self.set_status(500)
self.set_header("Cache-Control", "no-store")
self.set_header("Content-Type", "application/json")
async def _submit_upstream(self, upstream_provider_name, spot, credentials) -> str | None:
"""Submit a spot to the named upstream provider. Returns None on success, or a warning message on failure."""
provider = self._find_provider(upstream_provider_name, spot.activities)
if not provider:
if spot.activities:
return f"No enabled provider named '{upstream_provider_name}' supports upstream submission for {', '.join(spot.activities)} spots."
return f"No enabled provider named '{upstream_provider_name}' supports upstream submission for spots with no activity."
try:
# Submit spot to the upstream provider. Run in a separate thread otherwise this blocks the whole web server
# for everyone!
await IOLoop.current().run_in_executor(None, provider.submit_spot, spot, credentials)
logger.info(f"Spot of {spot.dx_call} by {spot.de_call} submitted upstream to {upstream_provider_name}.")
# Trigger a re-poll after 3 second so the spot appears quickly. (Submitting to a cluster node is slower than
# this, but we get data as a live stream from cluster anyway, so force_poll does nothing in that case. This
# is really just for the HTTP providers when we submit a spot to them)
threading.Timer(3.0, provider.force_poll).start()
return None
except NotImplementedError as e:
return str(e)
except SpotSubmissionError as e:
logger.warning(f"Upstream submission to {upstream_provider_name} was not accepted: {e}")
return f"Upstream submission to {upstream_provider_name} failed: {e}"
except Exception:
logger.exception(f"Failed to submit spot upstream to {upstream_provider_name}")
return f"Upstream submission to {upstream_provider_name} failed."
def _find_provider(self, provider_name, activities) -> SpotProvider | None:
"""Find an enabled provider by name that can submit spots for at least one of the given activities. If there
are no activities, find one that can submit spots with no activity."""
for p in self._spot_providers:
if p.enabled and p.name == provider_name and any(p.can_submit_spot(a) for a in activities or [None]):
return p
return None
@staticmethod
def _is_valid_api_key(api_key):
"""Check whether the supplied API key is one of the ones allowed in config."""
# HMAC Compare Digest is a recommended thing for security reasons. If you just compare strings
# then the comparison returns at the first non-matching character, which means in theory you can
# use the time it takes to compare strings to figure out how much of the string you've got right.
# With the digest approach it's not the real strings being compared but generated digests, so
# it will take a constant amount of time regardless of how well the actual strings match.
return any(hmac.compare_digest(api_key.encode(), k.encode()) for k in API_KEYS)
@staticmethod
def _verify_recaptcha(token):
"""Verify a Google reCAPTCHA v2 token. Returns True if valid."""
try:
response = requests.post(
RECAPTCHA_VERIFY_URL,
data={"secret": RECAPTCHA_SECRET_KEY, "response": token},
timeout=(5, 10),
)
return response.ok and response.json().get("success", False)
except Exception:
logger.exception("reCAPTCHA verification request failed")
return False