#!/usr/bin/python3 import base64 import grp import ipaddress import json import os import pwd import re import secrets import socket import socketserver import struct import threading import time import urllib.error import urllib.request from http.server import BaseHTTPRequestHandler from pathlib import Path from urllib.parse import urlsplit from cryptography.hazmat.primitives import hashes from cryptography.hazmat.primitives.ciphers.aead import AESGCM from cryptography.hazmat.primitives.kdf.hkdf import HKDF BASE_URL = "{{DROPS_ORIGIN}}" ORIGIN = "{{RELAY_ORIGIN}}" MANAGEMENT_ORIGIN = "{{DROPS_ORIGIN}}" SOCKET_PATH = Path("{{MANAGER_CONTROL_SOCKET_PATH}}") CATALOG_PATH = Path("{{MANAGER_CATALOG_PATH}}") DOWNLOAD_ROOT = Path("{{MANAGER_DOWNLOAD_ROOT}}") ASSET_ROOT = Path("/usr/local/share/neodrop-manager") HTTP_SOCKET_PATH = Path("{{MANAGER_HTTP_SOCKET_PATH}}") MANAGEMENT_NETWORKS = ( ipaddress.ip_network("{{MANAGEMENT_CIDR_1}}"), ipaddress.ip_network("{{MANAGEMENT_CIDR_2}}"), ipaddress.ip_network("{{MANAGEMENT_CIDR_3}}"), ) CATALOG_KEY_PATH = Path( "{{CATALOG_KEY_PATH}}" ) API_TOKEN_PATH = Path( "{{MANAGER_API_TOKEN_PATH}}" ) CHUNK_SIZE = 4 * 1024 * 1024 TEXT_READ_MAX = 1024 * 1024 TEXT_EDIT_MAX = 256 * 1024 TEXT_EDIT_LINES_MAX = 2000 HKDF_SALT = b"{{NEODROP_CRYPTO_CONTEXT}}" CATALOG_HEADER = b"NEODROP-CATALOG\x00\x01" OBJECT_ID_RE = re.compile(r"^[A-Za-z0-9_-]{32}$") TOKEN_RE = re.compile(r"^[A-Za-z0-9_-]{43}$") READ_URL_RE = re.compile( r"^https://{{DROPS_HOST_REGEX}}/d/([A-Za-z0-9_-]{32})" r"(?:\?k=(image|audio|video|text|binary))?#d1\.([A-Za-z0-9_-]{43})$" ) PROFILES = {0: 50 * 1024 * 1024, 1: 1024 * 1024 * 1024} PROTOCOL_PAYLOAD_MAX = 12 * 1024 class ManagerError(Exception): pass class Response: def __init__(self, status, body, headers=None): self.status = status self.body = body self.headers = headers or {} def json(self): try: return json.loads(self.body.decode("utf-8")) except (UnicodeDecodeError, json.JSONDecodeError) as error: raise ManagerError("invalid JSON response") from error class HTTPTransport: def request(self, method, path, body=None, headers=None): if not path.startswith("/api/v1/") or "?" in path or "#" in path: raise ManagerError("invalid API path") request = urllib.request.Request( BASE_URL + path, data=body, method=method, headers=headers or {} ) try: with urllib.request.urlopen(request, timeout=30) as response: return Response(response.status, response.read(), dict(response.headers)) except urllib.error.HTTPError as error: error.read() raise ManagerError("NeoDrop request failed (%d)" % error.code) from error except urllib.error.URLError as error: raise ManagerError("NeoDrop request unavailable") from error def b64url(data): return base64.urlsafe_b64encode(data).rstrip(b"=").decode("ascii") def unb64url(value, expected=None): try: raw = base64.urlsafe_b64decode(value + "=" * ((4 - len(value) % 4) % 4)) except (ValueError, TypeError) as error: raise ManagerError("invalid base64 value") from error if expected is not None and len(raw) != expected: raise ManagerError("invalid base64 value") if b64url(raw) != value: raise ManagerError("invalid base64 value") return raw def derive(secret, info, length): return HKDF( algorithm=hashes.SHA256(), length=length, salt=HKDF_SALT, info=info.encode("ascii"), ).derive(secret) def make_aad(record_type, profile, object_id, size, count, index, length): return ( b"NEODRP01" + bytes((record_type, profile, 0, 0)) + object_id + struct.pack(">QIIII", size, CHUNK_SIZE, count, index, length) ) def make_iv(record_type, index): return struct.pack(">IQ", record_type, index) def is_text_type(content_type): return bool( re.match(r"^text/", content_type, re.IGNORECASE) or re.match( r"^application/(json|javascript|x-.*script)$", content_type, re.IGNORECASE, ) ) def content_kind(content_type, name=""): if re.match(r"^image/(png|jpeg|gif|webp|avif)$", content_type, re.IGNORECASE): return "image" if re.match(r"^audio/", content_type, re.IGNORECASE): return "audio" if re.match(r"^video/", content_type, re.IGNORECASE): return "video" extension = Path(name).suffix.lower() if extension in (".png", ".jpg", ".jpeg", ".gif", ".webp", ".avif"): return "image" if extension in (".flac", ".m4a", ".mp3", ".oga", ".ogg", ".opus", ".wav", ".wave"): return "audio" if extension in (".m4v", ".mkv", ".mov", ".mp4", ".ogv", ".webm"): return "video" if is_text_type(content_type): return "text" return "binary" def safe_filename(value): name = str(value or "drop.bin").replace("\\", "/").split("/")[-1] name = "".join(char for char in name if ord(char) >= 32 and ord(char) != 127) name = name.strip().strip(".")[:240] return name or "drop.bin" class EncryptedCatalog: def __init__(self, path, key): if len(key) != 32: raise ManagerError("catalog key must be exactly 32 bytes") self.path = Path(path) self.key = key def load(self): try: payload = self.path.read_bytes() except FileNotFoundError: return {} minimum = len(CATALOG_HEADER) + 12 + 16 if len(payload) < minimum or not payload.startswith(CATALOG_HEADER): raise ManagerError("invalid encrypted catalog") nonce_at = len(CATALOG_HEADER) nonce = payload[nonce_at:nonce_at + 12] try: plaintext = AESGCM(self.key).decrypt( nonce, payload[nonce_at + 12:], CATALOG_HEADER ) document = json.loads(plaintext.decode("utf-8")) except Exception as error: raise ManagerError("unable to decrypt catalog") from error if not isinstance(document, dict) or document.get("v") != 1: raise ManagerError("invalid catalog version") entries = document.get("entries") if not isinstance(entries, list): raise ManagerError("invalid catalog entries") result = {} for entry in entries: if not isinstance(entry, dict) or not OBJECT_ID_RE.fullmatch(entry.get("id", "")): raise ManagerError("invalid catalog entry") result[entry["id"]] = entry return result def save(self, entries): self.path.parent.mkdir(mode=0o700, parents=True, exist_ok=True) document = {"v": 1, "entries": sorted(entries.values(), key=lambda item: item["id"])} plaintext = json.dumps( document, separators=(",", ":"), sort_keys=True, ensure_ascii=True ).encode("utf-8") nonce = secrets.token_bytes(12) payload = CATALOG_HEADER + nonce + AESGCM(self.key).encrypt( nonce, plaintext, CATALOG_HEADER ) temporary = self.path.with_name( ".%s.%s.tmp" % (self.path.name, secrets.token_hex(8)) ) try: with open(temporary, "xb") as handle: os.fchmod(handle.fileno(), 0o600) handle.write(payload) handle.flush() os.fsync(handle.fileno()) os.replace(temporary, self.path) directory = os.open(self.path.parent, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)) try: os.fsync(directory) finally: os.close(directory) finally: try: temporary.unlink() except FileNotFoundError: pass class NeoDropManager: def __init__(self, catalog, api_token, transport=None, download_root=None): if not api_token or any(ord(char) < 33 or ord(char) > 126 for char in api_token): raise ManagerError("invalid manager API token") self.catalog = catalog self.api_token = api_token self.transport = transport or HTTPTransport() self.download_root = Path(download_root or DOWNLOAD_ROOT) self.edit_session = None self.lock = threading.RLock() def request(self, method, path, body=None, headers=None): response = self.transport.request(method, path, body=body, headers=headers or {}) if response.status < 200 or response.status >= 300: raise ManagerError("NeoDrop request failed (%d)" % response.status) return response def json_request(self, method, path, value=None, headers=None): body = None request_headers = dict(headers or {}) if value is not None: body = json.dumps(value, separators=(",", ":")).encode("utf-8") request_headers["Content-Type"] = "application/json" return self.request(method, path, body, request_headers).json() def inventory(self): value = self.json_request( "GET", "/api/v1/manager/objects", headers={"Authorization": "Bearer " + self.api_token}, ) rows = value.get("objects") if isinstance(value, dict) else value if not isinstance(rows, list): raise ManagerError("invalid manager inventory") result = {} for row in rows: if not isinstance(row, dict): raise ManagerError("invalid manager inventory") object_id = row.get("id") or row.get("objectId") or row.get("object_id") if not OBJECT_ID_RE.fullmatch(object_id or ""): raise ManagerError("invalid manager inventory") result[object_id] = row return result def resolve(self, selector, entries=None): if not selector or not re.fullmatch(r"[A-Za-z0-9_-]{1,32}", selector): raise ManagerError("invalid selector") entries = entries if entries is not None else self.catalog.load() matches = [entry for object_id, entry in entries.items() if object_id.startswith(selector)] if not matches: raise ManagerError("object not found") if len(matches) != 1: raise ManagerError("selector is ambiguous") return matches[0] def parse_read_url(self, read_url): match = READ_URL_RE.fullmatch(read_url) if not match: raise ManagerError("invalid NeoDrop read URL") object_id, _, encoded_secret = match.groups() secret = unb64url(encoded_secret, 32) derived_id = derive(secret, "object-id", 24) if b64url(derived_id) != object_id: raise ManagerError("read URL identifier mismatch") return object_id, secret, derived_id def manifest_and_metadata(self, read_url): object_id, secret, object_id_raw = self.parse_read_url(read_url) manifest = self.json_request("GET", "/api/v1/objects/" + object_id) try: version = manifest["v"] profile = int(manifest["profile"]) size = int(manifest["size"]) chunk_size = int(manifest["chunkSize"]) chunks = int(manifest["chunks"]) expires_at = int(manifest["expiresAt"]) metadata_cipher = unb64url(manifest["metadata"]) except (KeyError, TypeError, ValueError): raise ManagerError("invalid object manifest") from None expected_chunks = (size + CHUNK_SIZE - 1) // CHUNK_SIZE if size else 0 if ( version != 1 or profile not in PROFILES or size < 0 or size > PROFILES[profile] or chunk_size != CHUNK_SIZE or chunks != expected_chunks or len(metadata_cipher) < 16 or len(metadata_cipher) > 4112 ): raise ManagerError("invalid object manifest") metadata_length = len(metadata_cipher) - 16 try: plaintext = AESGCM(derive(secret, "metadata-key", 32)).decrypt( make_iv(2, 0), metadata_cipher, make_aad( 0, profile, object_id_raw, size, chunks, 0xFFFFFFFF, metadata_length ), ) metadata = json.loads(plaintext.decode("utf-8")) except Exception as error: raise ManagerError("metadata authentication failed") from error if ( not isinstance(metadata, dict) or metadata.get("v") != 1 or metadata.get("size") != size or not isinstance(metadata.get("name"), str) or not isinstance(metadata.get("type"), str) ): raise ManagerError("invalid authenticated metadata") normalized = { "v": 1, "profile": profile, "size": size, "chunkSize": chunk_size, "chunks": chunks, "expiresAt": expires_at, } return object_id, secret, object_id_raw, normalized, metadata def register(self, read_url, delete_capability): object_id, _, _, manifest, metadata = self.manifest_and_metadata(read_url) parts = delete_capability.split(".") if ( len(parts) != 2 or parts[0] != object_id or not TOKEN_RE.fullmatch(parts[1]) ): raise ManagerError("invalid delete capability") entries = self.catalog.load() if object_id in entries: existing = entries[object_id] if ( existing.get("read_url") == read_url and existing.get("delete_capability") == delete_capability ): return ["registered %s generation %d" % (object_id, existing["generation"])] raise ManagerError("object is already registered") entry = { "id": object_id, "read_url": read_url, "delete_capability": delete_capability, "metadata": metadata, "size": manifest["size"], "profile": manifest["profile"], "chunks": manifest["chunks"], "expires_at": manifest["expiresAt"], "generation": 1, } entries[object_id] = entry self.catalog.save(entries) return ["registered %s generation 1" % object_id] def browser_entry(self, entry): metadata = entry.get("metadata", {}) return { "id": entry["id"], "name": metadata.get("name", "drop.bin"), "type": metadata.get("type", "application/octet-stream"), "size": entry["size"], "profile": entry["profile"], "expiresAt": entry["expires_at"], "readUrl": entry["read_url"], } def browser_list(self): entries = self.catalog.load() return [ self.browser_entry(entry) for entry in sorted( entries.values(), key=lambda item: (item["expires_at"], item["id"]) ) ] def browser_register(self, read_url, delete_capability): self.register(read_url, delete_capability) object_id, _, _ = self.parse_read_url(read_url) return self.browser_entry(self.catalog.load()[object_id]) def browser_delete(self, object_id): if not OBJECT_ID_RE.fullmatch(object_id or ""): raise ManagerError("invalid object identifier") self.delete(object_id) def list_objects(self): entries = self.catalog.load() inventory = self.inventory() lines = [] for object_id in sorted(set(entries) | set(inventory)): if object_id in entries: entry = entries[object_id] lines.append( "%s managed gen=%d size=%d profile=%d name=%s" % ( object_id, entry["generation"], entry["size"], entry["profile"], entry["metadata"]["name"], ) ) else: row = inventory[object_id] lines.append( "%s unmanaged size=%s profile=%s" % (object_id, row.get("size", "?"), row.get("profile", "?")) ) return lines or ["no objects"] def show(self, selector): entry = self.resolve(selector) metadata = json.dumps( entry["metadata"], separators=(",", ":"), sort_keys=True, ensure_ascii=True ) return [ "id %s" % entry["id"], "generation %d" % entry["generation"], "profile %d" % entry["profile"], "size %d" % entry["size"], "chunks %d" % entry["chunks"], "expires %d" % entry["expires_at"], "metadata %s" % metadata, "url %s" % entry["read_url"], ] def prepare_content(self, entry, maximum=None): object_id, secret, object_id_raw, manifest, metadata = self.manifest_and_metadata( entry["read_url"] ) if object_id != entry["id"]: raise ManagerError("catalog identifier mismatch") if maximum is not None and manifest["size"] > maximum: raise ManagerError("object is too large") key = AESGCM(derive(secret, "content-key", 32)) return object_id, object_id_raw, manifest, metadata, key def decrypted_chunks(self, prepared): object_id, object_id_raw, manifest, _, key = prepared for index in range(manifest["chunks"]): plain_length = min(CHUNK_SIZE, manifest["size"] - index * CHUNK_SIZE) cipher = self.request( "GET", "/api/v1/objects/%s/chunks/%d" % (object_id, index) ).body if len(cipher) != plain_length + 16: raise ManagerError("invalid encrypted chunk length") try: plain = key.decrypt( make_iv(1, index), cipher, make_aad( 1, manifest["profile"], object_id_raw, manifest["size"], manifest["chunks"], index, plain_length, ), ) except Exception as error: raise ManagerError("content authentication failed") from error yield plain def decrypt_content(self, entry, maximum=None): prepared = self.prepare_content(entry, maximum) manifest = prepared[2] metadata = prepared[3] parts = list(self.decrypted_chunks(prepared)) return b"".join(parts), metadata, manifest def read_text(self, selector, edit_limit=False): entry = self.resolve(selector) maximum = TEXT_EDIT_MAX if edit_limit else TEXT_READ_MAX content, metadata, _ = self.decrypt_content(entry, maximum) if not is_text_type(metadata["type"]): raise ManagerError("object is not text") try: text = content.decode("utf-8") except UnicodeDecodeError as error: raise ManagerError("text is not valid UTF-8") from error return entry, text, metadata def read_command(self, selector): _, text, _ = self.read_text(selector) return text.split("\n") def unique_download_path(self, name): self.download_root.mkdir(mode=0o700, parents=True, exist_ok=True) safe = safe_filename(name) candidate = self.download_root / safe stem = candidate.stem suffix = candidate.suffix number = 1 while candidate.exists(): candidate = self.download_root / ("%s (%d)%s" % (stem, number, suffix)) number += 1 return candidate def download(self, selector): entry = self.resolve(selector) prepared = self.prepare_content(entry) metadata = prepared[3] destination = self.unique_download_path(metadata["name"]) temporary = destination.with_name(".%s.%s.tmp" % (destination.name, secrets.token_hex(8))) try: with open(temporary, "xb") as handle: os.fchmod(handle.fileno(), 0o600) for chunk in self.decrypted_chunks(prepared): handle.write(chunk) handle.flush() os.fsync(handle.fileno()) os.replace(temporary, destination) finally: try: temporary.unlink() except FileNotFoundError: pass return ["downloaded %s" % destination, "url %s" % entry["read_url"]] def delete_remote(self, entry): object_id, token = entry["delete_capability"].split(".", 1) if object_id != entry["id"] or not TOKEN_RE.fullmatch(token): raise ManagerError("invalid stored delete capability") self.json_request( "DELETE", "/api/v1/objects/" + object_id, headers={"Authorization": "Bearer " + token, "Origin": ORIGIN}, ) def delete(self, selector): entries = self.catalog.load() entry = self.resolve(selector, entries) self.delete_remote(entry) del entries[entry["id"]] self.catalog.save(entries) if self.edit_session and self.edit_session["entry"]["id"] == entry["id"]: self.edit_session = None return ["deleted %s" % entry["id"]] def numbered_context(self, index): lines = self.edit_session["lines"] if not lines: return ["(empty)"] index = max(0, min(index, len(lines) - 1)) first = max(0, index - 2) last = min(len(lines), index + 3) return ["%4d | %s" % (position + 1, lines[position]) for position in range(first, last)] def edit(self, selector): entry, text, _ = self.read_text(selector, edit_limit=True) lines = text.split("\n") trailing_newline = bool(lines and lines[-1] == "") if trailing_newline: lines.pop() if len(lines) > TEXT_EDIT_LINES_MAX: raise ManagerError("text has too many lines") self.edit_session = { "entry": entry, "lines": lines, "trailing_newline": trailing_newline, } return ["editing %s (%d lines)" % (entry["id"], len(lines))] + self.numbered_context(0) def require_edit(self): if self.edit_session is None: raise ManagerError("no edit session") def mutate_line(self, operation, number, text=None): self.require_edit() if number < 1: raise ManagerError("line number out of range") lines = self.edit_session["lines"] original = list(lines) if operation == "line": if number > len(lines): raise ManagerError("line number out of range") lines[number - 1] = text focus = number - 1 elif operation == "insert": if number > len(lines) + 1: raise ManagerError("line number out of range") lines.insert(number - 1, text) focus = number - 1 else: if number > len(lines): raise ManagerError("line number out of range") del lines[number - 1] focus = max(0, number - 2) if len(lines) > TEXT_EDIT_LINES_MAX: lines[:] = original raise ManagerError("text has too many lines") if len(self.edit_bytes()) > TEXT_EDIT_MAX: lines[:] = original raise ManagerError("edited text is too large") return self.numbered_context(focus) def edit_bytes(self): text = "\n".join(self.edit_session["lines"]) if self.edit_session["trailing_newline"]: text += "\n" return text.encode("utf-8") def upload(self, content, metadata, profile, generation): if profile not in PROFILES or len(content) > PROFILES[profile]: raise ManagerError("content exceeds profile limit") secret = secrets.token_bytes(32) object_id_raw = derive(secret, "object-id", 24) object_id = b64url(object_id_raw) chunks = (len(content) + CHUNK_SIZE - 1) // CHUNK_SIZE if content else 0 session = self.json_request( "POST", "/api/v1/uploads", { "v": 1, "id": object_id, "profile": profile, "size": len(content), "chunkSize": CHUNK_SIZE, "chunks": chunks, }, {"Origin": ORIGIN}, ) upload_id = session.get("uploadId") if isinstance(session, dict) else None upload_token = session.get("uploadToken") if isinstance(session, dict) else None if not re.fullmatch(r"[A-Za-z0-9_-]{22}", upload_id or "") or not TOKEN_RE.fullmatch( upload_token or "" ): raise ManagerError("invalid upload session") stored_metadata = dict(metadata) stored_metadata.update({"v": 1, "size": len(content), "name": safe_filename(metadata["name"])}) stored_metadata["lastModified"] = int(time.time() * 1000) metadata_plain = json.dumps( stored_metadata, separators=(",", ":"), ensure_ascii=True ).encode("utf-8") if len(metadata_plain) > 4096: raise ManagerError("metadata is too large") metadata_cipher = AESGCM(derive(secret, "metadata-key", 32)).encrypt( make_iv(2, 0), metadata_plain, make_aad( 0, profile, object_id_raw, len(content), chunks, 0xFFFFFFFF, len(metadata_plain), ), ) authorization = "Bearer " + upload_token binary_headers = { "Authorization": authorization, "Content-Type": "application/octet-stream", "Origin": ORIGIN, } self.request( "PUT", "/api/v1/uploads/%s/metadata" % upload_id, metadata_cipher, binary_headers, ) content_key = AESGCM(derive(secret, "content-key", 32)) for index in range(chunks): plain = content[index * CHUNK_SIZE:(index + 1) * CHUNK_SIZE] cipher = content_key.encrypt( make_iv(1, index), plain, make_aad( 1, profile, object_id_raw, len(content), chunks, index, len(plain), ), ) self.request( "PUT", "/api/v1/uploads/%s/chunks/%d" % (upload_id, index), cipher, binary_headers, ) commit = self.json_request( "POST", "/api/v1/uploads/%s/commit" % upload_id, {}, {"Authorization": authorization, "Origin": ORIGIN}, ) try: expires_at = int(commit["expiresAt"]) except (KeyError, TypeError, ValueError): raise ManagerError("invalid upload commit response") from None read_url = "%s/d/%s?k=%s#d1.%s" % ( BASE_URL, object_id, content_kind(stored_metadata["type"], stored_metadata["name"]), b64url(secret), ) return { "id": object_id, "read_url": read_url, "delete_capability": object_id + "." + upload_token, "metadata": stored_metadata, "size": len(content), "profile": profile, "chunks": chunks, "expires_at": expires_at, "generation": generation, } def save_edit(self): self.require_edit() old_entry = self.edit_session["entry"] content = self.edit_bytes() new_entry = self.upload( content, old_entry["metadata"], old_entry["profile"], old_entry["generation"] + 1, ) entries = self.catalog.load() current = entries.get(old_entry["id"]) if current is None or current.get("generation") != old_entry["generation"]: raise ManagerError("catalog changed during edit") del entries[old_entry["id"]] entries[new_entry["id"]] = new_entry self.catalog.save(entries) self.edit_session = None lines = [ "saved %s generation %d" % (new_entry["id"], new_entry["generation"]), "url %s" % new_entry["read_url"], ] try: self.delete_remote(old_entry) except ManagerError: lines.append("warning old object deletion failed") return lines def put_text(self, name, text): content = text.encode("utf-8") if len(content) > TEXT_EDIT_MAX: raise ManagerError("text is too large") entry = self.upload( content, {"v": 1, "name": safe_filename(name), "type": "text/plain", "size": len(content)}, 0, 1, ) entries = self.catalog.load() entries[entry["id"]] = entry self.catalog.save(entries) return [ "created %s generation 1" % entry["id"], "url %s" % entry["read_url"], ] def cancel_edit(self): self.require_edit() self.edit_session = None return ["edit cancelled"] def handle_command(self, command): with self.lock: command = command.rstrip("\r") if command == "help": return [ "help | status | list | register READ_URL DELETE_CAPABILITY", "put-text NAME TEXT", "show SELECTOR | read SELECTOR | download SELECTOR | delete SELECTOR", "edit SELECTOR | line N TEXT | insert N TEXT | delete-line N | save | cancel", ] if command == "status": return [ "catalog=%d inventory=%d edit=%s" % ( len(self.catalog.load()), len(self.inventory()), "active" if self.edit_session else "none", ) ] if command == "list": return self.list_objects() match = re.fullmatch(r"register\s+(\S+)\s+(\S+)", command) if match: return self.register(match.group(1), match.group(2)) match = re.fullmatch(r"put-text\s+(\S+)\s+(.*)", command) if match: return self.put_text(match.group(1), match.group(2)) match = re.fullmatch(r"(show|read|download|delete|edit)\s+(\S+)", command) if match: action, selector = match.groups() return { "show": self.show, "read": self.read_command, "download": self.download, "delete": self.delete, "edit": self.edit, }[action](selector) match = re.fullmatch(r"(line|insert)\s+(\d+)\s(.*)", command) if match: return self.mutate_line(match.group(1), int(match.group(2)), match.group(3)) match = re.fullmatch(r"delete-line\s+(\d+)", command) if match: return self.mutate_line("delete-line", int(match.group(1))) if command == "save": return self.save_edit() if command == "cancel": return self.cancel_edit() raise ManagerError("unknown command; use help") def protocol_text(value): return str(value).replace("\x00", "").replace("\r", " ").replace("\n", " ") def protocol_records(record_type, value): if record_type not in ("LINE", "ERROR"): raise ValueError("invalid protocol record type") payload = protocol_text(value).encode("utf-8") prefix = (record_type + "\t").encode("ascii") if not payload: yield prefix + b"\n" return while payload: cut = min(PROTOCOL_PAYLOAD_MAX, len(payload)) if cut < len(payload): while cut > 0 and payload[cut] & 0xC0 == 0x80: cut -= 1 yield prefix + payload[:cut] + b"\n" payload = payload[cut:] def handle_client(client, manager): buffer = b"" try: while True: data = client.recv(4096) if not data: return buffer += data if len(buffer) > 1024 * 1024: raise ManagerError("command is too large") while b"\n" in buffer: line, buffer = buffer.split(b"\n", 1) try: text = line.decode("utf-8") prefix, command = text.split("\t", 1) if prefix != "COMMAND": raise ManagerError("expected COMMAND record") for output in manager.handle_command(command): for record in protocol_records("LINE", output): client.sendall(record) except (ManagerError, UnicodeDecodeError, ValueError) as error: for record in protocol_records("ERROR", protocol_text(error)[:300]): client.sendall(record) except (ManagerError, OSError): pass finally: client.close() class ManagementHandler(BaseHTTPRequestHandler): protocol_version = "HTTP/1.1" server_version = "NeoDrop-Manager" sys_version = "" def log_message(self, fmt, *args): # Request bodies contain capabilities; keep this endpoint deliberately quiet. return def client_ip(self): try: return ipaddress.ip_address(self.headers.get("X-Real-IP", "")) except ValueError: return None def source_allowed(self): if not self.server.peer_allowed(self.connection): return False address = self.client_ip() return ( self.headers.get("Host") == "{{DROPS_HOST}}" and address is not None and any(address in network for network in MANAGEMENT_NETWORKS) ) def api_allowed(self): return self.source_allowed() and self.headers.get("Origin") == MANAGEMENT_ORIGIN def send_body(self, status, content_type, body): self.send_response(status) self.send_header("Content-Type", content_type) self.send_header("Content-Length", str(len(body))) self.send_header("Cache-Control", "private, no-store, max-age=0") self.send_header("Pragma", "no-cache") self.send_header("Referrer-Policy", "no-referrer") self.send_header("X-Content-Type-Options", "nosniff") self.send_header("X-Frame-Options", "DENY") self.send_header( "Content-Security-Policy", "default-src 'none'; script-src 'self'; style-src 'self'; " "connect-src 'self'; frame-ancestors 'none'; " "base-uri 'none'; form-action 'none'", ) self.end_headers() if self.command != "HEAD": self.wfile.write(body) def send_json(self, status, value): body = json.dumps(value, separators=(",", ":"), ensure_ascii=True).encode("utf-8") + b"\n" self.send_body(status, "application/json", body) def discard_request_body(self): try: length = int(self.headers.get("Content-Length", "0")) except ValueError: return if 0 < length <= 4096: self.rfile.read(length) def forbidden(self): self.discard_request_body() self.close_connection = True self.send_json(403, {"error": "forbidden"}) def unavailable(self): self.discard_request_body() self.close_connection = True self.send_json(404, {"error": "unavailable"}) def read_api_json(self, maximum): if self.headers.get("Content-Type") != "application/json": self.discard_request_body() self.close_connection = True return None try: length = int(self.headers.get("Content-Length", "-1")) except ValueError: self.close_connection = True return None if length < 0 or length > maximum: self.discard_request_body() self.close_connection = True return None try: return json.loads(self.rfile.read(length).decode("utf-8")) except (UnicodeDecodeError, json.JSONDecodeError): self.close_connection = True return None def static_asset(self, path): assets = { "/manage/": ("manage.html", "text/html; charset=utf-8"), "/manage/style.css": ("manage.css", "text/css; charset=utf-8"), "/manage/manage.js": ("manage.js", "application/javascript"), "/manage/neodrop-crypto.js": ("neodrop-crypto.js", "application/javascript"), } asset = assets.get(path) if not asset: return self.unavailable() try: body = (ASSET_ROOT / asset[0]).read_bytes() except OSError: return self.unavailable() self.send_body(200, asset[1], body) def do_HEAD(self): return self.do_GET() def do_GET(self): parsed = urlsplit(self.path) if self.path != parsed.path or not self.source_allowed(): return self.forbidden() return self.static_asset(parsed.path) def do_POST(self): parsed = urlsplit(self.path) if self.path != parsed.path or not self.api_allowed(): return self.forbidden() if parsed.path == "/manage/api/v1/list": data = self.read_api_json(64) if data != {}: return self.send_json(400, {"error": "invalid request"}) try: with self.server.manager.lock: rows = self.server.manager.browser_list() except ManagerError: return self.send_json(503, {"error": "unavailable"}) return self.send_json(200, {"drops": rows}) if parsed.path == "/manage/api/v1/register": data = self.read_api_json(2048) if not isinstance(data, dict) or set(data) != {"readUrl", "deleteCapability"}: return self.send_json(400, {"error": "invalid request"}) if not isinstance(data["readUrl"], str) or not isinstance(data["deleteCapability"], str): return self.send_json(400, {"error": "invalid request"}) try: with self.server.manager.lock: row = self.server.manager.browser_register( data["readUrl"], data["deleteCapability"] ) except ManagerError: return self.send_json(400, {"error": "invalid request"}) return self.send_json(201, {"drop": row}) return self.unavailable() def do_DELETE(self): parsed = urlsplit(self.path) match = re.fullmatch(r"/manage/api/v1/drops/([A-Za-z0-9_-]{32})", parsed.path) if self.path != parsed.path or not self.api_allowed(): return self.forbidden() if not match or self.headers.get("Content-Length", "0") != "0": return self.unavailable() try: with self.server.manager.lock: self.server.manager.browser_delete(match.group(1)) except ManagerError: return self.send_json(404, {"error": "unavailable"}) self.send_json(200, {"deleted": True}) def reject_method(self): parsed = urlsplit(self.path) if self.path != parsed.path or not self.api_allowed(): return self.forbidden() return self.unavailable() do_OPTIONS = reject_method do_PATCH = reject_method do_PUT = reject_method def linux_peer_uid(connection): option = getattr(socket, "SO_PEERCRED", None) if option is None: raise OSError("SO_PEERCRED is unavailable") credentials = connection.getsockopt( socket.SOL_SOCKET, option, struct.calcsize("3i") ) _, uid, _ = struct.unpack("3i", credentials) return uid class ManagementHTTPServer(socketserver.ThreadingMixIn, socketserver.UnixStreamServer): daemon_threads = True request_queue_size = 32 def __init__(self, address, manager, expected_peer_uid, peer_uid=linux_peer_uid): super().__init__(address, ManagementHandler) self.manager = manager self.expected_peer_uid = expected_peer_uid self.peer_uid = peer_uid def peer_allowed(self, connection): try: return self.peer_uid(connection) == self.expected_peer_uid except (OSError, struct.error, ValueError): return False def load_credentials(): key = CATALOG_KEY_PATH.read_bytes() if len(key) != 32: raise ManagerError("catalog key credential must be exactly 32 raw bytes") try: token = API_TOKEN_PATH.read_bytes().strip().decode("ascii") except UnicodeDecodeError as error: raise ManagerError("invalid manager API token credential") from error return key, token def socket_server(manager): SOCKET_PATH.parent.mkdir(mode=0o750, parents=True, exist_ok=True) try: SOCKET_PATH.unlink() except FileNotFoundError: pass server = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) server.bind(str(SOCKET_PATH)) os.chmod(SOCKET_PATH, 0o660) os.chown(SOCKET_PATH, -1, grp.getgrnam("neodrop-control").gr_gid) server.listen(8) while True: client, _ = server.accept() threading.Thread(target=handle_client, args=(client, manager), daemon=True).start() def main(): os.umask(0o077) key, api_token = load_credentials() manager = NeoDropManager(EncryptedCatalog(CATALOG_PATH, key), api_token) try: HTTP_SOCKET_PATH.unlink() except FileNotFoundError: pass nginx_uid = pwd.getpwnam("www-data").pw_uid http_server = ManagementHTTPServer( str(HTTP_SOCKET_PATH), manager, expected_peer_uid=nginx_uid ) os.chmod(HTTP_SOCKET_PATH, 0o660) threading.Thread(target=http_server.serve_forever, daemon=True).start() socket_server(manager) if __name__ == "__main__": main()