Source-Tag: v1.0.3 Manifest-SHA256: b7b09bf48f9e092ed8ad82f99c0319d6c3837d0a1549bd0cc4a374276c0f3896
1144 lines
42 KiB
Python
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()
|