Files
neorelay/files/neodrop-manager/neodrop_manager_agent.py
T
NIA Sanitized Release Publisher 7141326ea0 Sanitized public release v1.0.3
Source-Tag: v1.0.3
Manifest-SHA256: b7b09bf48f9e092ed8ad82f99c0319d6c3837d0a1549bd0cc4a374276c0f3896
2026-07-20 22:35:36 +00:00

1144 lines
42 KiB
Python

#!/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()