Restructuring as a Python package.
This commit is contained in:
@@ -0,0 +1,95 @@
|
||||
import asyncio
|
||||
import os
|
||||
import unicodedata
|
||||
from pathlib import Path
|
||||
|
||||
from pathvalidate import sanitize_filepath
|
||||
|
||||
from .asynclink import AsyncLink
|
||||
from .lrucache import LRUCache
|
||||
|
||||
ROOT = Path(os.environ.get("STORAGE", Path.cwd()))
|
||||
|
||||
def sanitize_filename(filename):
|
||||
filename = unicodedata.normalize("NFC", filename)
|
||||
filename = sanitize_filepath(filename)
|
||||
filename = filename.replace("/", "-")
|
||||
return filename
|
||||
|
||||
class File:
|
||||
def __init__(self, filename):
|
||||
self.path = ROOT / filename
|
||||
self.fd = None
|
||||
self.writable = False
|
||||
|
||||
def open_ro(self):
|
||||
self.close()
|
||||
self.fd = os.open(self.path, os.O_RDONLY)
|
||||
|
||||
def open_rw(self):
|
||||
self.close()
|
||||
self.path.parent.mkdir(parents=True, exist_ok=True)
|
||||
self.fd = os.open(self.path, os.O_RDWR | os.O_CREAT)
|
||||
self.writable = True
|
||||
|
||||
def write(self, pos, buffer, *, file_size=None):
|
||||
if not self.writable:
|
||||
# Create/open file
|
||||
self.open_rw()
|
||||
if file_size is not None:
|
||||
os.ftruncate(self.fd, file_size)
|
||||
os.lseek(self.fd, pos, os.SEEK_SET)
|
||||
os.write(self.fd, buffer)
|
||||
|
||||
def __getitem__(self, slice):
|
||||
if self.fd is None:
|
||||
self.open_ro()
|
||||
os.lseek(self.fd, slice.start, os.SEEK_SET)
|
||||
l = slice.stop - slice.start
|
||||
data = os.read(self.fd, l)
|
||||
if len(data) < l: raise EOFError("Error reading requested range")
|
||||
return data
|
||||
|
||||
def close(self):
|
||||
if self.fd is not None:
|
||||
os.close(self.fd)
|
||||
self.fd = self.writable = None
|
||||
|
||||
def __del__(self):
|
||||
self.close()
|
||||
|
||||
|
||||
class FileServer:
|
||||
|
||||
async def start(self):
|
||||
self.alink = AsyncLink()
|
||||
self.worker = asyncio.get_event_loop().run_in_executor(None, self.worker_thread, self.alink.to_sync)
|
||||
self.cache = LRUCache(File, capacity=10, maxage=5.0)
|
||||
|
||||
async def stop(self):
|
||||
await self.alink.stop()
|
||||
await self.worker
|
||||
|
||||
def worker_thread(self, slink):
|
||||
try:
|
||||
for req in slink:
|
||||
with req as (command, *args):
|
||||
if command == "upload":
|
||||
req.set_result(self.upload(*args))
|
||||
elif command == "download":
|
||||
req.set_result(self.download(*args))
|
||||
else:
|
||||
raise NotImplementedError(f"Unhandled {command=} {args}")
|
||||
finally:
|
||||
self.cache.close()
|
||||
|
||||
def upload(self, name, pos, data, file_size):
|
||||
name = sanitize_filename(name)
|
||||
f = self.cache[name]
|
||||
f.write(pos, data, file_size=file_size)
|
||||
return len(data)
|
||||
|
||||
def download(self, name, start, end):
|
||||
name = sanitize_filename(name)
|
||||
f = self.cache[name]
|
||||
return f[start: end]
|
||||
Reference in New Issue
Block a user