anonymousland-synapse/synapse/rest/media/v1/media_storage.py

275 lines
9.0 KiB
Python
Raw Normal View History

2018-01-08 12:07:30 -05:00
# -*- coding: utf-8 -*-
# Copyright 2018 New Vecotr Ltd
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
from twisted.internet import defer, threads
from twisted.protocols.basic import FileSender
from ._base import Responder
2018-01-17 11:56:23 -05:00
from synapse.util.file_consumer import BackgroundFileConsumer
2018-01-08 12:07:30 -05:00
from synapse.util.logcontext import make_deferred_yieldable
import contextlib
import os
import logging
import shutil
2018-01-12 09:44:02 -05:00
import sys
2018-01-08 12:07:30 -05:00
2018-01-17 11:56:23 -05:00
2018-01-08 12:07:30 -05:00
logger = logging.getLogger(__name__)
class MediaStorage(object):
"""Responsible for storing/fetching files from local sources.
2018-01-12 06:15:31 -05:00
Args:
local_media_directory (str): Base path where we store media on disk
filepaths (MediaFilePaths)
storage_providers ([StorageProvider]): List of StorageProvider that are
used to fetch and store files.
2018-01-08 12:07:30 -05:00
"""
2018-01-08 12:19:55 -05:00
def __init__(self, local_media_directory, filepaths, storage_providers):
2018-01-08 12:07:30 -05:00
self.local_media_directory = local_media_directory
self.filepaths = filepaths
2018-01-08 12:19:55 -05:00
self.storage_providers = storage_providers
2018-01-08 12:07:30 -05:00
@defer.inlineCallbacks
def store_file(self, source, file_info):
"""Write `source` to the on disk media store, and also any other
configured storage providers
Args:
source: A file like object that should be written
file_info (FileInfo): Info about the file to store
Returns:
Deferred[str]: the file path written to in the primary media store
"""
path = self._file_info_to_path(file_info)
fname = os.path.join(self.local_media_directory, path)
dirname = os.path.dirname(fname)
if not os.path.exists(dirname):
os.makedirs(dirname)
# Write to the main repository
yield make_deferred_yieldable(threads.deferToThread(
_write_file_synchronously, source, fname,
))
# Tell the storage providers about the new file. They'll decide
# if they should upload it and whether to do so synchronously
# or not.
for provider in self.storage_providers:
yield provider.store_file(path, file_info)
2018-01-08 12:07:30 -05:00
defer.returnValue(fname)
@contextlib.contextmanager
def store_into_file(self, file_info):
"""Context manager used to get a file like object to write into, as
described by file_info.
2018-01-12 06:15:31 -05:00
Actually yields a 3-tuple (file, fname, finish_cb), where file is a file
like object that can be written to, fname is the absolute path of file
on disk, and finish_cb is a function that returns a Deferred.
fname can be used to read the contents from after upload, e.g. to
generate thumbnails.
finish_cb must be called and waited on after the file has been
successfully been written to. Should not be called if there was an
error.
2018-01-08 12:07:30 -05:00
Args:
file_info (FileInfo): Info about the file to store
Example:
with media_storage.store_into_file(info) as (f, fname, finish_cb):
# .. write into f ...
yield finish_cb()
"""
path = self._file_info_to_path(file_info)
fname = os.path.join(self.local_media_directory, path)
dirname = os.path.dirname(fname)
if not os.path.exists(dirname):
os.makedirs(dirname)
finished_called = [False]
2018-01-08 12:19:55 -05:00
@defer.inlineCallbacks
2018-01-08 12:07:30 -05:00
def finish():
2018-01-08 12:19:55 -05:00
for provider in self.storage_providers:
yield provider.store_file(path, file_info)
2018-01-08 12:07:30 -05:00
finished_called[0] = True
try:
with open(fname, "wb") as f:
yield f, fname, finish
2018-01-12 09:44:02 -05:00
except Exception:
t, v, tb = sys.exc_info()
2018-01-08 12:07:30 -05:00
try:
os.remove(fname)
except Exception:
pass
2018-01-12 09:44:02 -05:00
raise t, v, tb
2018-01-08 12:07:30 -05:00
if not finished_called:
2018-01-12 06:15:31 -05:00
raise Exception("Finished callback not called")
2018-01-08 12:07:30 -05:00
@defer.inlineCallbacks
def fetch_media(self, file_info):
"""Attempts to fetch media described by file_info from the local cache
and configured storage providers.
Args:
file_info (FileInfo)
Returns:
2018-01-12 06:15:31 -05:00
Deferred[Responder|None]: Returns a Responder if the file was found,
2018-01-08 12:07:30 -05:00
otherwise None.
"""
path = self._file_info_to_path(file_info)
local_path = os.path.join(self.local_media_directory, path)
if os.path.exists(local_path):
defer.returnValue(FileResponder(open(local_path, "rb")))
2018-01-08 12:19:55 -05:00
for provider in self.storage_providers:
res = yield provider.fetch(path, file_info)
if res:
defer.returnValue(res)
2018-01-08 12:07:30 -05:00
defer.returnValue(None)
@defer.inlineCallbacks
def ensure_media_is_in_local_cache(self, file_info):
"""Ensures that the given file is in the local cache. Attempts to
download it from storage providers if it isn't.
Args:
file_info (FileInfo)
Returns:
Deferred[str]: Full path to local file
"""
path = self._file_info_to_path(file_info)
local_path = os.path.join(self.local_media_directory, path)
if os.path.exists(local_path):
defer.returnValue(local_path)
2018-01-17 11:56:23 -05:00
dirname = os.path.dirname(local_path)
if not os.path.exists(dirname):
os.makedirs(dirname)
for provider in self.storage_providers:
res = yield provider.fetch(path, file_info)
if res:
with res:
2018-01-17 11:56:23 -05:00
consumer = BackgroundFileConsumer(open(local_path, "w"))
yield res.write_to_consumer(consumer)
yield consumer.wait()
defer.returnValue(local_path)
raise Exception("file could not be found")
2018-01-08 12:07:30 -05:00
def _file_info_to_path(self, file_info):
"""Converts file_info into a relative path.
2018-01-12 06:15:31 -05:00
The path is suitable for storing files under a directory, e.g. used to
store files on local FS under the base media repository directory.
Args:
file_info (FileInfo)
Returns:
str
2018-01-08 12:07:30 -05:00
"""
if file_info.url_cache:
if file_info.thumbnail:
return self.filepaths.url_cache_thumbnail_rel(
media_id=file_info.file_id,
width=file_info.thumbnail_width,
height=file_info.thumbnail_height,
content_type=file_info.thumbnail_type,
method=file_info.thumbnail_method,
)
2018-01-08 12:07:30 -05:00
return self.filepaths.url_cache_filepath_rel(file_info.file_id)
if file_info.server_name:
if file_info.thumbnail:
return self.filepaths.remote_media_thumbnail_rel(
server_name=file_info.server_name,
file_id=file_info.file_id,
width=file_info.thumbnail_width,
height=file_info.thumbnail_height,
content_type=file_info.thumbnail_type,
method=file_info.thumbnail_method
)
return self.filepaths.remote_media_filepath_rel(
file_info.server_name, file_info.file_id,
)
if file_info.thumbnail:
return self.filepaths.local_media_thumbnail_rel(
media_id=file_info.file_id,
width=file_info.thumbnail_width,
height=file_info.thumbnail_height,
content_type=file_info.thumbnail_type,
method=file_info.thumbnail_method
)
return self.filepaths.local_media_filepath_rel(
file_info.file_id,
)
def _write_file_synchronously(source, fname):
"""Write `source` to the path `fname` synchronously. Should be called
from a thread.
Args:
source: A file like object to be written
fname (str): Path to write to
"""
dirname = os.path.dirname(fname)
if not os.path.exists(dirname):
os.makedirs(dirname)
source.seek(0) # Ensure we read from the start of the file
with open(fname, "wb") as f:
shutil.copyfileobj(source, f)
class FileResponder(Responder):
"""Wraps an open file that can be sent to a request.
Args:
open_file (file): A file like object to be streamed ot the client,
is closed when finished streaming.
"""
def __init__(self, open_file):
self.open_file = open_file
def write_to_consumer(self, consumer):
2018-01-17 11:56:23 -05:00
return FileSender().beginFileTransfer(self.open_file, consumer)
2018-01-08 12:07:30 -05:00
2018-01-12 08:32:03 -05:00
def __exit__(self, exc_type, exc_val, exc_tb):
2018-01-08 12:07:30 -05:00
self.open_file.close()