mirror of
https://github.com/Novartis/cellxgene-gateway.git
synced 2026-09-25 06:58:11 +08:00
Compare commits
24
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a4443cb57e | ||
|
|
c92fe18ee7 | ||
|
|
54aacab276 | ||
|
|
7dd10f1f3d | ||
|
|
8f676f28d0 | ||
|
|
a74576ade5 | ||
|
|
d747860118 | ||
|
|
79fef57010 | ||
|
|
cd9c0a3671 | ||
|
|
55b268125f | ||
|
|
4df58f9ceb | ||
|
|
8b4565e745 | ||
|
|
c8056991b0 | ||
|
|
1d1d8b4e59 | ||
|
|
a35c6b9b1e | ||
|
|
f260180a76 | ||
|
|
58ae41fe0c | ||
|
|
0c5adc9fef | ||
|
|
34ac73ba01 | ||
|
|
d16a97906c | ||
|
|
3e5accad65 | ||
|
|
903d25763f | ||
|
|
08c546f40a | ||
|
|
b9f4d35812 |
@@ -40,6 +40,7 @@ jobs:
|
|||||||
eval "$(conda shell.bash hook)"
|
eval "$(conda shell.bash hook)"
|
||||||
conda activate cellxgene-gateway
|
conda activate cellxgene-gateway
|
||||||
python setup.py install
|
python setup.py install
|
||||||
|
pip install -r requirements-test.txt
|
||||||
|
|
||||||
- name: Run tests
|
- name: Run tests
|
||||||
run: |
|
run: |
|
||||||
@@ -52,9 +53,32 @@ jobs:
|
|||||||
eval "$(conda shell.bash hook)"
|
eval "$(conda shell.bash hook)"
|
||||||
conda activate cellxgene-gateway
|
conda activate cellxgene-gateway
|
||||||
coverage report --fail-under 41
|
coverage report --fail-under 41
|
||||||
|
coverage report > coverage.txt
|
||||||
|
coverage html -i
|
||||||
coverage xml -i
|
coverage xml -i
|
||||||
|
|
||||||
|
- name: Upload coverage HTML report
|
||||||
|
uses: actions/upload-artifact@v4
|
||||||
|
with:
|
||||||
|
name: coverage-html
|
||||||
|
path: htmlcov/
|
||||||
|
retention-days: 30
|
||||||
|
|
||||||
|
- name: Upload coverage xml
|
||||||
|
uses: actions/upload-artifact@v4
|
||||||
|
with:
|
||||||
|
name: coverage-xml
|
||||||
|
path: coverage.xml
|
||||||
|
retention-days: 30
|
||||||
|
|
||||||
|
- name: Upload coverage summary
|
||||||
|
uses: actions/upload-artifact@v4
|
||||||
|
with:
|
||||||
|
name: coverage-summary
|
||||||
|
path: coverage.txt
|
||||||
|
retention-days: 30
|
||||||
- name: "Upload coverage to Codecov"
|
- name: "Upload coverage to Codecov"
|
||||||
|
if: ${{ github.event_name == 'push' || (github.event_name == 'pull_request' && github.event.pull_request.head.repo.full_name == github.repository) }}
|
||||||
uses: codecov/codecov-action@v1
|
uses: codecov/codecov-action@v1
|
||||||
with:
|
with:
|
||||||
token: ${{ secrets.CODECOV_TOKEN }}
|
token: ${{ secrets.CODECOV_TOKEN }}
|
||||||
|
|||||||
@@ -1,3 +1,19 @@
|
|||||||
|
# 0.4.2
|
||||||
|
|
||||||
|
* update package name
|
||||||
|
|
||||||
|
# 0.4.1
|
||||||
|
|
||||||
|
* Fix UnicodeDecodeError when viewing compressed datasets
|
||||||
|
* Fix WSGI server initialization by extracting data source setup
|
||||||
|
* Fix AttributeError by storing ItemSource objects in default_item_source
|
||||||
|
* Delay itemsource initialization until first request is served
|
||||||
|
* Set default_item_source and start pruner thread
|
||||||
|
* Added start scripts for flask, gunicorn and uwsgi
|
||||||
|
* Made pruner a daemon thread
|
||||||
|
* Updated start scripts to run in subshells
|
||||||
|
* Fixed bug in status.json
|
||||||
|
|
||||||
# 0.4.0
|
# 0.4.0
|
||||||
|
|
||||||
* Removed dependency on flask-api
|
* Removed dependency on flask-api
|
||||||
|
|||||||
@@ -115,6 +115,21 @@ docker run -it --rm \
|
|||||||
-p 8080:8080 \
|
-p 8080:8080 \
|
||||||
cellxgene-gateway
|
cellxgene-gateway
|
||||||
```
|
```
|
||||||
|
## Running cellxgene gateway with start scripts
|
||||||
|
|
||||||
|
For your convenience, we provide start scripts for flask, gunicorn and uwsgi.
|
||||||
|
|
||||||
|
First, set up a .env
|
||||||
|
```bash
|
||||||
|
cp env_example .env
|
||||||
|
# edit .env
|
||||||
|
open .env
|
||||||
|
```
|
||||||
|
|
||||||
|
Then run the scripts in a subshell
|
||||||
|
```bash
|
||||||
|
( ./start_flask.sh )
|
||||||
|
```
|
||||||
|
|
||||||
# Customization
|
# Customization
|
||||||
|
|
||||||
|
|||||||
@@ -7,4 +7,4 @@
|
|||||||
# OR CONDITIONS OF ANY KIND, either express or implied. See the License for
|
# OR CONDITIONS OF ANY KIND, either express or implied. See the License for
|
||||||
# the specific language governing permissions and limitations under the License.
|
# the specific language governing permissions and limitations under the License.
|
||||||
|
|
||||||
__version__ = "0.4.0"
|
__version__ = "0.4.2"
|
||||||
|
|||||||
@@ -148,7 +148,7 @@ class CacheEntry:
|
|||||||
headers = {}
|
headers = {}
|
||||||
copy_headers = [
|
copy_headers = [
|
||||||
"accept",
|
"accept",
|
||||||
"accept-encoding",
|
# "accept-encoding" - removed: let requests library handle compression/decompression
|
||||||
"accept-language",
|
"accept-language",
|
||||||
"cache-control",
|
"cache-control",
|
||||||
"connection",
|
"connection",
|
||||||
@@ -168,9 +168,8 @@ class CacheEntry:
|
|||||||
headers[h] = request.headers[h]
|
headers[h] = request.headers[h]
|
||||||
|
|
||||||
full_path = self.cellxgene_basepath() + subpath + querystring()
|
full_path = self.cellxgene_basepath() + subpath + querystring()
|
||||||
|
cellxgene_response = None
|
||||||
try:
|
try:
|
||||||
cellxgene_response = None
|
|
||||||
if request.method in ["GET", "HEAD", "OPTIONS"]:
|
if request.method in ["GET", "HEAD", "OPTIONS"]:
|
||||||
cellxgene_response = get(full_path, headers=headers)
|
cellxgene_response = get(full_path, headers=headers)
|
||||||
elif request.method == "PUT":
|
elif request.method == "PUT":
|
||||||
|
|||||||
+133
-51
@@ -40,6 +40,12 @@ app = Flask(__name__)
|
|||||||
item_sources = []
|
item_sources = []
|
||||||
default_item_source = None
|
default_item_source = None
|
||||||
|
|
||||||
|
# Guard for lazy initialization so tests can import this module without
|
||||||
|
# triggering environment-dependent side effects. initialize_data_sources()
|
||||||
|
# will set this to True when it has run.
|
||||||
|
data_sources_initialized = False
|
||||||
|
data_sources_init_lock = Lock()
|
||||||
|
|
||||||
|
|
||||||
def _force_https(app):
|
def _force_https(app):
|
||||||
def wrapper(environ, start_response):
|
def wrapper(environ, start_response):
|
||||||
@@ -75,9 +81,87 @@ if (
|
|||||||
x_prefix=env.proxy_fix_prefix,
|
x_prefix=env.proxy_fix_prefix,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# WSGI middleware to ensure data sources are initialized before the first
|
||||||
|
# WSGI request is handled. This guarantees initialization works under
|
||||||
|
# Gunicorn/uWSGI (which import the module but don't call main()). The
|
||||||
|
# initialize_data_sources() function is idempotent-protected by
|
||||||
|
# data_sources_initialized and data_sources_init_lock.
|
||||||
|
def _init_on_first_wsgi_request(wsgi_app):
|
||||||
|
def middleware(environ, start_response):
|
||||||
|
global data_sources_initialized
|
||||||
|
if not data_sources_initialized:
|
||||||
|
with data_sources_init_lock:
|
||||||
|
if not app.extensions.get("cellxgene_gateway", {}).get("launchtime"):
|
||||||
|
app.extensions.setdefault("cellxgene_gateway", {})[
|
||||||
|
"launchtime"
|
||||||
|
] = current_time_stamp()
|
||||||
|
|
||||||
|
if not data_sources_initialized:
|
||||||
|
initialize_data_sources()
|
||||||
|
|
||||||
|
env.validate()
|
||||||
|
if not item_sources or not len(item_sources):
|
||||||
|
raise Exception(
|
||||||
|
"No data sources specified for Cellxgene Gateway"
|
||||||
|
)
|
||||||
|
|
||||||
|
global default_item_source
|
||||||
|
if default_item_source is None:
|
||||||
|
default_item_source = item_sources[0]
|
||||||
|
|
||||||
|
data_sources_initialized = True
|
||||||
|
return wsgi_app(environ, start_response)
|
||||||
|
|
||||||
|
return middleware
|
||||||
|
|
||||||
|
|
||||||
|
# Wrap the WSGI app so Gunicorn/uWSGI will trigger initialization when the
|
||||||
|
# first request comes in. Tests that need initialization can call
|
||||||
|
# initialize_data_sources() directly.
|
||||||
|
app.wsgi_app = _init_on_first_wsgi_request(app.wsgi_app)
|
||||||
|
|
||||||
cache = BackendCache()
|
cache = BackendCache()
|
||||||
|
|
||||||
|
|
||||||
|
# Initialize data sources - this is defined later in the file but called here
|
||||||
|
# to ensure initialization happens when WSGI servers (Gunicorn) import the module
|
||||||
|
def initialize_data_sources():
|
||||||
|
"""Initialize data sources from environment variables.
|
||||||
|
Called at module import time for WSGI server compatibility (Gunicorn).
|
||||||
|
Uses a guard flag to prevent double initialization within a process."""
|
||||||
|
global default_item_source
|
||||||
|
|
||||||
|
logging.basicConfig(
|
||||||
|
level=env.log_level,
|
||||||
|
format="%(asctime)s:%(name)s:%(levelname)s:%(message)s",
|
||||||
|
)
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
cellxgene_data = os.environ.get("CELLXGENE_DATA", None)
|
||||||
|
cellxgene_bucket = os.environ.get("CELLXGENE_BUCKET", None)
|
||||||
|
|
||||||
|
if cellxgene_bucket is not None:
|
||||||
|
from cellxgene_gateway.items.s3.s3item_source import S3ItemSource
|
||||||
|
|
||||||
|
s3_source = S3ItemSource(cellxgene_bucket, name="s3")
|
||||||
|
item_sources.append(s3_source)
|
||||||
|
default_item_source = s3_source
|
||||||
|
logger.info("Initialized S3 data source")
|
||||||
|
logger.debug(f"S3 bucket: {cellxgene_bucket}")
|
||||||
|
if cellxgene_data is not None:
|
||||||
|
from cellxgene_gateway.items.file.fileitem_source import FileItemSource
|
||||||
|
|
||||||
|
file_source = FileItemSource(cellxgene_data, name="local")
|
||||||
|
item_sources.append(file_source)
|
||||||
|
default_item_source = file_source
|
||||||
|
logger.info("Initialized local file data source")
|
||||||
|
logger.debug(f"Data directory: {cellxgene_data}")
|
||||||
|
if len(item_sources) == 0:
|
||||||
|
raise Exception("Please specify CELLXGENE_DATA or CELLXGENE_BUCKET")
|
||||||
|
flask_util.include_source_in_url = len(item_sources) > 1
|
||||||
|
|
||||||
|
|
||||||
@app.errorhandler(CellxgeneException)
|
@app.errorhandler(CellxgeneException)
|
||||||
def handle_invalid_usage(error):
|
def handle_invalid_usage(error):
|
||||||
message = f"{error.http_status} Error : {error.message}"
|
message = f"{error.http_status} Error : {error.message}"
|
||||||
@@ -169,7 +253,7 @@ entry_lock = Lock()
|
|||||||
|
|
||||||
|
|
||||||
def matching_source(source_name):
|
def matching_source(source_name):
|
||||||
if source_name is None:
|
if source_name is None and default_item_source is not None:
|
||||||
source_name = default_item_source.name
|
source_name = default_item_source.name
|
||||||
matching = [i for i in item_sources if i.name == source_name]
|
matching = [i for i in item_sources if i.name == source_name]
|
||||||
if len(matching) != 1:
|
if len(matching) != 1:
|
||||||
@@ -216,6 +300,11 @@ def do_view(path, source_name=None):
|
|||||||
raise CellxgeneException("User not authorized to access this data", 403)
|
raise CellxgeneException("User not authorized to access this data", 403)
|
||||||
elif match.status == CacheEntryStatus.error:
|
elif match.status == CacheEntryStatus.error:
|
||||||
raise ProcessException.from_cache_entry(match)
|
raise ProcessException.from_cache_entry(match)
|
||||||
|
else:
|
||||||
|
raise CellxgeneException(
|
||||||
|
f"Unexpected cache entry status {match.status} for key {match.key.descriptor}",
|
||||||
|
500,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@app.route("/cache_status", methods=["GET"])
|
@app.route("/cache_status", methods=["GET"])
|
||||||
@@ -229,28 +318,40 @@ def do_GET_status():
|
|||||||
|
|
||||||
@app.route("/cache_status.json", methods=["GET"])
|
@app.route("/cache_status.json", methods=["GET"])
|
||||||
def do_GET_status_json():
|
def do_GET_status_json():
|
||||||
|
def map_entry(entry):
|
||||||
|
dataset = entry.key.h5ad_item.descriptor
|
||||||
|
annotation_file = entry.key.annotation_descriptor
|
||||||
|
return {
|
||||||
|
"dataset": dataset,
|
||||||
|
"annotation_file": annotation_file,
|
||||||
|
"launchtime": entry.launchtime,
|
||||||
|
"last_access": entry.timestamp,
|
||||||
|
"status": entry.status.name,
|
||||||
|
}
|
||||||
|
|
||||||
return json.dumps(
|
return json.dumps(
|
||||||
{
|
{
|
||||||
"launchtime": app.launchtime,
|
"launchtime": app.extensions.get("cellxgene_gateway", {}).get("launchtime"),
|
||||||
"entry_list": [
|
"entry_list": [map_entry(entry) for entry in cache.entry_list],
|
||||||
{
|
|
||||||
"dataset": entry.key.dataset,
|
|
||||||
"annotation_file": entry.key.annotation_file,
|
|
||||||
"launchtime": entry.launchtime,
|
|
||||||
"last_access": entry.timestamp,
|
|
||||||
"status": entry.status,
|
|
||||||
}
|
|
||||||
for entry in cache.entry_list
|
|
||||||
],
|
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@app.route("/relaunch/<path:path>", methods=["GET"])
|
def get_cache_key(path):
|
||||||
def do_relaunch(path):
|
if request.args.get("source_name"):
|
||||||
source_name = request.args.get("source_name") or default_item_source.name
|
source_name = request.args.get("source_name")
|
||||||
|
elif default_item_source:
|
||||||
|
source_name = default_item_source.name
|
||||||
|
else:
|
||||||
|
source_name = None
|
||||||
source = matching_source(source_name)
|
source = matching_source(source_name)
|
||||||
key = CacheKey.for_lookup(source, source.lookup(path))
|
key = CacheKey.for_lookup(source, source.lookup(path))
|
||||||
|
return key
|
||||||
|
|
||||||
|
|
||||||
|
@app.route("/relaunch/<path:path>", methods=["GET"])
|
||||||
|
def do_relaunch(path):
|
||||||
|
key = get_cache_key(path)
|
||||||
match = cache.check_entry(key)
|
match = cache.check_entry(key)
|
||||||
if not match is None:
|
if not match is None:
|
||||||
match.terminate()
|
match.terminate()
|
||||||
@@ -262,9 +363,7 @@ def do_relaunch(path):
|
|||||||
|
|
||||||
@app.route("/terminate/<path:path>", methods=["GET"])
|
@app.route("/terminate/<path:path>", methods=["GET"])
|
||||||
def do_terminate(path):
|
def do_terminate(path):
|
||||||
source_name = request.args.get("source_name") or default_item_source.name
|
key = get_cache_key(path)
|
||||||
source = matching_source(source_name)
|
|
||||||
key = CacheKey.for_lookup(source, source.lookup(path))
|
|
||||||
match = cache.check_entry(key)
|
match = cache.check_entry(key)
|
||||||
if not match is None:
|
if not match is None:
|
||||||
match.terminate()
|
match.terminate()
|
||||||
@@ -277,46 +376,29 @@ def ip_address():
|
|||||||
return set_no_cache(resp)
|
return set_no_cache(resp)
|
||||||
|
|
||||||
|
|
||||||
def launch():
|
def start_pruner_thread():
|
||||||
env.validate()
|
|
||||||
if not item_sources or not len(item_sources):
|
|
||||||
raise Exception("No data sources specified for Cellxgene Gateway")
|
|
||||||
|
|
||||||
global default_item_source
|
|
||||||
if default_item_source is None:
|
|
||||||
default_item_source = item_sources[0]
|
|
||||||
|
|
||||||
pruner = PruneProcessCache(cache)
|
pruner = PruneProcessCache(cache)
|
||||||
|
# Run the pruner as a daemon thread so it won't block interpreter
|
||||||
background_thread = Thread(target=pruner)
|
# shutdown (for example when Ctrl-C is used in the main thread).
|
||||||
|
# This avoids "Exception ignored in: <module 'threading'...>" at exit.
|
||||||
|
background_thread = Thread(target=pruner, daemon=True)
|
||||||
background_thread.start()
|
background_thread.start()
|
||||||
|
|
||||||
app.launchtime = current_time_stamp()
|
|
||||||
|
def launch():
|
||||||
|
start_pruner_thread()
|
||||||
|
|
||||||
|
app.extensions.setdefault("cellxgene_gateway", {})[
|
||||||
|
"launchtime"
|
||||||
|
] = current_time_stamp()
|
||||||
app.run(host="0.0.0.0", port=env.gateway_port, debug=False)
|
app.run(host="0.0.0.0", port=env.gateway_port, debug=False)
|
||||||
|
|
||||||
|
|
||||||
|
app.extensions.setdefault("cellxgene_gateway", {})["launchtime"] = None
|
||||||
|
|
||||||
|
|
||||||
def main():
|
def main():
|
||||||
logging.basicConfig(
|
"""CLI entry point for Flask development server."""
|
||||||
level=env.log_level,
|
|
||||||
format="%(asctime)s:%(name)s:%(levelname)s:%(message)s",
|
|
||||||
)
|
|
||||||
cellxgene_data = os.environ.get("CELLXGENE_DATA", None)
|
|
||||||
cellxgene_bucket = os.environ.get("CELLXGENE_BUCKET", None)
|
|
||||||
|
|
||||||
if cellxgene_bucket is not None:
|
|
||||||
from cellxgene_gateway.items.s3.s3item_source import S3ItemSource
|
|
||||||
|
|
||||||
item_sources.append(S3ItemSource(cellxgene_bucket, name="s3"))
|
|
||||||
default_item_source = "s3"
|
|
||||||
if cellxgene_data is not None:
|
|
||||||
from cellxgene_gateway.items.file.fileitem_source import FileItemSource
|
|
||||||
|
|
||||||
item_sources.append(FileItemSource(cellxgene_data, name="local"))
|
|
||||||
default_item_source = "local"
|
|
||||||
if len(item_sources) == 0:
|
|
||||||
raise Exception("Please specify CELLXGENE_DATA or CELLXGENE_BUCKET")
|
|
||||||
flask_util.include_source_in_url = len(item_sources) > 1
|
|
||||||
|
|
||||||
launch()
|
launch()
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,3 @@
|
|||||||
|
export CELLXGENE_LOCATION=$(pwd)/.venv/bin/cellxgene
|
||||||
|
export CELLXGENE_DATA=../cellxgene_data
|
||||||
|
export GATEWAY_IP=127.0.0.1
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
freezegun~=1.5.5
|
||||||
@@ -1,6 +0,0 @@
|
|||||||
export CELLXGENE_LOCATION=$(pwd)/.cellxgene-gateway/bin/cellxgene
|
|
||||||
export CELLXGENE_DATA=../cellxgene_data
|
|
||||||
export GATEWAY_IP=127.0.0.1
|
|
||||||
|
|
||||||
#Once these are set, you run like a normal Flask app
|
|
||||||
cellxgene-gateway
|
|
||||||
@@ -38,7 +38,7 @@ install_reqs = parse_requirements()
|
|||||||
|
|
||||||
setup(
|
setup(
|
||||||
# mandatory
|
# mandatory
|
||||||
name="cellxgene-gateway",
|
name="cellxgene_gateway",
|
||||||
# mandatory
|
# mandatory
|
||||||
version=get_version("cellxgene_gateway/__init__.py"),
|
version=get_version("cellxgene_gateway/__init__.py"),
|
||||||
# mandatory
|
# mandatory
|
||||||
|
|||||||
Executable
+41
@@ -0,0 +1,41 @@
|
|||||||
|
#!/bin/bash
|
||||||
|
|
||||||
|
# start_gunicorn.sh - Start Cellxgene Gateway with Gunicorn
|
||||||
|
#
|
||||||
|
# PREREQUISITES:
|
||||||
|
# - Gunicorn installed (included with cellxgene 1.3.0, or: pip install gunicorn)
|
||||||
|
# - Virtual environment activated
|
||||||
|
# - .env file with CELLXGENE_LOCATION and CELLXGENE_DATA (or CELLXGENE_BUCKET)
|
||||||
|
#
|
||||||
|
# USAGE:
|
||||||
|
# ./start_gunicorn.sh
|
||||||
|
|
||||||
|
# Exit on error
|
||||||
|
set -e
|
||||||
|
|
||||||
|
# Get the directory where this script is located
|
||||||
|
SCRIPT_DIR="$( cd "$( dirname "${BASH_SOURCE[0]}" )" && pwd )"
|
||||||
|
|
||||||
|
|
||||||
|
# Source environment variables
|
||||||
|
echo "Loading environment variables..."
|
||||||
|
if [ -f "$SCRIPT_DIR/.env" ]; then
|
||||||
|
source "$SCRIPT_DIR/.env"
|
||||||
|
else
|
||||||
|
echo "Error: .env file not found at $SCRIPT_DIR/.env"
|
||||||
|
echo "Please create it with required environment variables"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
# Verify required environment variables
|
||||||
|
if [ -z "$CELLXGENE_LOCATION" ]; then
|
||||||
|
echo "Error: CELLXGENE_LOCATION not set"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
if [ -z "$CELLXGENE_DATA" ] && [ -z "$CELLXGENE_BUCKET" ]; then
|
||||||
|
echo "Error: Either CELLXGENE_DATA or CELLXGENE_BUCKET must be set"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
exec cellxgene-gateway
|
||||||
Executable
+91
@@ -0,0 +1,91 @@
|
|||||||
|
#!/bin/bash
|
||||||
|
|
||||||
|
# start_gunicorn.sh - Start Cellxgene Gateway with Gunicorn
|
||||||
|
#
|
||||||
|
# PREREQUISITES:
|
||||||
|
# - Gunicorn installed (included with cellxgene 1.3.0, or: pip install gunicorn)
|
||||||
|
# - Virtual environment activated
|
||||||
|
# - .env file with CELLXGENE_LOCATION and CELLXGENE_DATA (or CELLXGENE_BUCKET)
|
||||||
|
#
|
||||||
|
# USAGE:
|
||||||
|
# ./start_gunicorn.sh
|
||||||
|
|
||||||
|
# Exit on error
|
||||||
|
set -e
|
||||||
|
|
||||||
|
# Get the directory where this script is located
|
||||||
|
SCRIPT_DIR="$( cd "$( dirname "${BASH_SOURCE[0]}" )" && pwd )"
|
||||||
|
|
||||||
|
# Source environment variables
|
||||||
|
echo "Loading environment variables..."
|
||||||
|
if [ -f "$SCRIPT_DIR/.env" ]; then
|
||||||
|
source "$SCRIPT_DIR/.env"
|
||||||
|
else
|
||||||
|
echo "Error: .env file not found at $SCRIPT_DIR/.env"
|
||||||
|
echo "Please create it with required environment variables"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
# Verify required environment variables
|
||||||
|
if [ -z "$CELLXGENE_LOCATION" ]; then
|
||||||
|
echo "Error: CELLXGENE_LOCATION not set"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
if [ -z "$CELLXGENE_DATA" ] && [ -z "$CELLXGENE_BUCKET" ]; then
|
||||||
|
echo "Error: Either CELLXGENE_DATA or CELLXGENE_BUCKET must be set"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
# Gunicorn configuration
|
||||||
|
# WARNING: Multi-worker mode has cache synchronization issues (see plans/002-shared-cache-implementation.md)
|
||||||
|
# Each worker maintains its own in-memory cache, causing 404s for static assets when different
|
||||||
|
# workers handle requests for the same dataset. Use GUNICORN_WORKERS=1 until shared cache is implemented.
|
||||||
|
WORKERS=${GUNICORN_WORKERS:-1}
|
||||||
|
BIND=${GATEWAY_IP:-0.0.0.0}:${GATEWAY_PORT:-5005}
|
||||||
|
TIMEOUT=${GUNICORN_TIMEOUT:-120}
|
||||||
|
WORKER_CLASS=${GUNICORN_WORKER_CLASS:-sync}
|
||||||
|
KEEPALIVE=${GUNICORN_KEEPALIVE:-5}
|
||||||
|
LOG_LEVEL=${GUNICORN_LOG_LEVEL:-info}
|
||||||
|
|
||||||
|
# Production optimization: enable backed mode to reduce memory usage
|
||||||
|
export GATEWAY_ENABLE_BACKED_MODE=${GATEWAY_ENABLE_BACKED_MODE:-true}
|
||||||
|
|
||||||
|
# Check if gunicorn is installed
|
||||||
|
if ! command -v gunicorn &> /dev/null; then
|
||||||
|
echo "Error: gunicorn not found. Install with: pip install gunicorn"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
# Display configuration
|
||||||
|
echo "Starting Cellxgene Gateway with Gunicorn..."
|
||||||
|
echo "Configuration:"
|
||||||
|
echo " Data source: ${CELLXGENE_DATA:-$CELLXGENE_BUCKET}"
|
||||||
|
echo " Binding to: $BIND"
|
||||||
|
echo " Workers: $WORKERS"
|
||||||
|
echo " Worker class: $WORKER_CLASS"
|
||||||
|
echo " Timeout: ${TIMEOUT}s"
|
||||||
|
echo " Keepalive: ${KEEPALIVE}s"
|
||||||
|
echo " Log level: $LOG_LEVEL"
|
||||||
|
echo " Backed mode: ${GATEWAY_ENABLE_BACKED_MODE}"
|
||||||
|
echo ""
|
||||||
|
|
||||||
|
cd "$SCRIPT_DIR"
|
||||||
|
|
||||||
|
# Start Gunicorn with optimized settings
|
||||||
|
# Additional options you can add via environment variables:
|
||||||
|
# - GUNICORN_MAX_REQUESTS: Restart worker after N requests (prevents memory leaks)
|
||||||
|
# - GUNICORN_MAX_REQUESTS_JITTER: Add randomness to max-requests
|
||||||
|
exec gunicorn cellxgene_gateway.gateway:app \
|
||||||
|
--workers "$WORKERS" \
|
||||||
|
--worker-class "$WORKER_CLASS" \
|
||||||
|
--bind "$BIND" \
|
||||||
|
--timeout "$TIMEOUT" \
|
||||||
|
--keep-alive "$KEEPALIVE" \
|
||||||
|
--access-logfile - \
|
||||||
|
--error-logfile - \
|
||||||
|
--log-level "$LOG_LEVEL" \
|
||||||
|
--preload \
|
||||||
|
${GUNICORN_MAX_REQUESTS:+--max-requests "$GUNICORN_MAX_REQUESTS"} \
|
||||||
|
${GUNICORN_MAX_REQUESTS_JITTER:+--max-requests-jitter "$GUNICORN_MAX_REQUESTS_JITTER"} \
|
||||||
|
"$@"
|
||||||
Executable
+89
@@ -0,0 +1,89 @@
|
|||||||
|
#!/bin/bash
|
||||||
|
|
||||||
|
# start_uwsgi.sh - Start Cellxgene Gateway with uWSGI
|
||||||
|
#
|
||||||
|
# PREREQUISITES:
|
||||||
|
# - uWSGI installed (pip install uwsgi)
|
||||||
|
# - Virtual environment activated
|
||||||
|
# - .env file with CELLXGENE_LOCATION and CELLXGENE_DATA (or CELLXGENE_BUCKET)
|
||||||
|
#
|
||||||
|
# USAGE:
|
||||||
|
# ./start_uwsgi.sh
|
||||||
|
|
||||||
|
# Exit on error
|
||||||
|
set -e
|
||||||
|
|
||||||
|
# Get the directory where this script is located
|
||||||
|
SCRIPT_DIR="$( cd "$( dirname "${BASH_SOURCE[0]}" )" && pwd )"
|
||||||
|
|
||||||
|
# Source environment variables
|
||||||
|
echo "Loading environment variables..."
|
||||||
|
if [ -f "$SCRIPT_DIR/.env" ]; then
|
||||||
|
source "$SCRIPT_DIR/.env"
|
||||||
|
else
|
||||||
|
echo "Error: .env file not found at $SCRIPT_DIR/.env"
|
||||||
|
echo "Please create it with required environment variables"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
# Verify required environment variables
|
||||||
|
if [ -z "$CELLXGENE_LOCATION" ]; then
|
||||||
|
echo "Error: CELLXGENE_LOCATION not set"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
if [ -z "$CELLXGENE_DATA" ] && [ -z "$CELLXGENE_BUCKET" ]; then
|
||||||
|
echo "Error: Either CELLXGENE_DATA or CELLXGENE_BUCKET must be set"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
# uWSGI configuration
|
||||||
|
# WARNING: Multi-worker mode has cache synchronization issues (see plans/002-shared-cache-implementation.md)
|
||||||
|
# Each worker maintains its own in-memory cache, causing 404s for static assets when different
|
||||||
|
# workers handle requests for the same dataset. Use UWSGI_WORKERS=1 until shared cache is implemented.
|
||||||
|
WORKERS=${UWSGI_WORKERS:-1}
|
||||||
|
HOST=${GATEWAY_IP:-0.0.0.0}
|
||||||
|
PORT=${GATEWAY_PORT:-5005}
|
||||||
|
TIMEOUT=${UWSGI_TIMEOUT:-120}
|
||||||
|
THREADS=${UWSGI_THREADS:-1}
|
||||||
|
|
||||||
|
# Production optimization: enable backed mode to reduce memory usage
|
||||||
|
export GATEWAY_ENABLE_BACKED_MODE=${GATEWAY_ENABLE_BACKED_MODE:-true}
|
||||||
|
|
||||||
|
# Check if uwsgi is installed
|
||||||
|
if ! command -v uwsgi &> /dev/null; then
|
||||||
|
echo "Error: uwsgi not found. Install with: pip install uwsgi"
|
||||||
|
exit 1
|
||||||
|
else
|
||||||
|
# Display configuration
|
||||||
|
echo "Starting Cellxgene Gateway with uWSGI..."
|
||||||
|
echo "Configuration:"
|
||||||
|
echo " Data source: ${CELLXGENE_DATA:-$CELLXGENE_BUCKET}"
|
||||||
|
echo " Binding to: $HOST:$PORT"
|
||||||
|
echo " Workers: $WORKERS"
|
||||||
|
echo " Threads: $THREADS"
|
||||||
|
echo " Timeout: ${TIMEOUT}s"
|
||||||
|
echo " Backed mode: ${GATEWAY_ENABLE_BACKED_MODE}"
|
||||||
|
echo ""
|
||||||
|
|
||||||
|
cd "$SCRIPT_DIR"
|
||||||
|
|
||||||
|
# Start uWSGI with optimized settings
|
||||||
|
# Additional options you can add via environment variables:
|
||||||
|
# - UWSGI_MAX_REQUESTS: Restart worker after N requests (prevents memory leaks)
|
||||||
|
exec uwsgi \
|
||||||
|
--http "$HOST:$PORT" \
|
||||||
|
--module cellxgene_gateway.gateway:app \
|
||||||
|
--workers "$WORKERS" \
|
||||||
|
--threads "$THREADS" \
|
||||||
|
--harakiri "$TIMEOUT" \
|
||||||
|
--master \
|
||||||
|
--enable-threads \
|
||||||
|
--single-interpreter \
|
||||||
|
--need-app \
|
||||||
|
--die-on-term \
|
||||||
|
--log-x-forwarded-for \
|
||||||
|
${UWSGI_MAX_REQUESTS:+--max-requests "$UWSGI_MAX_REQUESTS"} \
|
||||||
|
"$@"
|
||||||
|
fi
|
||||||
|
|
||||||
@@ -1,13 +1,24 @@
|
|||||||
import unittest
|
import unittest
|
||||||
|
import os
|
||||||
|
import shutil
|
||||||
|
import tempfile
|
||||||
|
|
||||||
from unittest.mock import MagicMock, Mock, patch
|
from unittest.mock import MagicMock, Mock, patch
|
||||||
|
|
||||||
from cellxgene_gateway.gateway import app
|
|
||||||
from cellxgene_gateway.items.item import ItemType
|
from cellxgene_gateway.items.item import ItemType
|
||||||
from cellxgene_gateway.items.s3.s3item import S3Item
|
from cellxgene_gateway.items.s3.s3item import S3Item
|
||||||
from cellxgene_gateway.items.s3.s3item_source import S3ItemSource
|
from cellxgene_gateway.items.s3.s3item_source import S3ItemSource
|
||||||
|
from cellxgene_gateway.gateway import app
|
||||||
|
|
||||||
|
|
||||||
class TestScanDirectory(unittest.TestCase):
|
class TestScanDirectory(unittest.TestCase):
|
||||||
|
def setUp(self):
|
||||||
|
|
||||||
|
self.app = app
|
||||||
|
|
||||||
|
def tearDown(self):
|
||||||
|
pass
|
||||||
|
|
||||||
@patch("s3fs.S3FileSystem")
|
@patch("s3fs.S3FileSystem")
|
||||||
def test_GIVEN_invalid_bucket_THEN_throws_error(self, s3func):
|
def test_GIVEN_invalid_bucket_THEN_throws_error(self, s3func):
|
||||||
class S3Mock:
|
class S3Mock:
|
||||||
@@ -26,6 +37,7 @@ class TestScanDirectory(unittest.TestCase):
|
|||||||
|
|
||||||
@patch("s3fs.S3FileSystem")
|
@patch("s3fs.S3FileSystem")
|
||||||
def test__GIVEN_multilevel_bucket_THEN_properly_recurses_suburls(self, s3func):
|
def test__GIVEN_multilevel_bucket_THEN_properly_recurses_suburls(self, s3func):
|
||||||
|
|
||||||
class S3Mock:
|
class S3Mock:
|
||||||
def exists(path):
|
def exists(path):
|
||||||
if path in [
|
if path in [
|
||||||
@@ -82,7 +94,7 @@ class TestScanDirectory(unittest.TestCase):
|
|||||||
|
|
||||||
s3func.return_value = S3Mock
|
s3func.return_value = S3Mock
|
||||||
source = S3ItemSource("my-bucket")
|
source = S3ItemSource("my-bucket")
|
||||||
with app.test_request_context(query_string="refresh=true") as test_context:
|
with self.app.test_request_context(query_string="refresh=true") as test_context:
|
||||||
tree = source.scan_directory()
|
tree = source.scan_directory()
|
||||||
|
|
||||||
def s3item_compare(i1, i2, msg=""):
|
def s3item_compare(i1, i2, msg=""):
|
||||||
|
|||||||
+466
-2
@@ -1,7 +1,13 @@
|
|||||||
import unittest
|
import unittest
|
||||||
from unittest.mock import MagicMock, patch
|
import pytest
|
||||||
|
from http import HTTPStatus
|
||||||
|
from unittest.mock import MagicMock, Mock, patch, call
|
||||||
|
from freezegun import freeze_time
|
||||||
|
|
||||||
from cellxgene_gateway.backend_cache import is_port_in_use
|
from cellxgene_gateway.backend_cache import BackendCache, is_port_in_use
|
||||||
|
from cellxgene_gateway.cache_entry import CacheEntry, CacheEntryStatus
|
||||||
|
from cellxgene_gateway.cache_key import CacheKey
|
||||||
|
from cellxgene_gateway.cellxgene_exception import CellxgeneException
|
||||||
|
|
||||||
|
|
||||||
class TestIsPortInUse(unittest.TestCase):
|
class TestIsPortInUse(unittest.TestCase):
|
||||||
@@ -26,3 +32,461 @@ class TestIsPortInUse(unittest.TestCase):
|
|||||||
self.assertTrue(connectMock.connect_ex.calledOnceWith("a"))
|
self.assertTrue(connectMock.connect_ex.calledOnceWith("a"))
|
||||||
self.assertTrue(socketMock.calledOnceWith("a"))
|
self.assertTrue(socketMock.calledOnceWith("a"))
|
||||||
self.assertEqual(is_port_in_use(123), False)
|
self.assertEqual(is_port_in_use(123), False)
|
||||||
|
|
||||||
|
|
||||||
|
class TestBackendCacheInit(unittest.TestCase):
|
||||||
|
def test_GIVEN_new_backend_cache_THEN_entry_list_is_empty(self):
|
||||||
|
"""BackendCache should initialize with an empty entry list."""
|
||||||
|
cache = BackendCache()
|
||||||
|
self.assertEqual(cache.entry_list, [])
|
||||||
|
|
||||||
|
|
||||||
|
class TestBackendCacheGetPorts(unittest.TestCase):
|
||||||
|
def test_GIVEN_empty_cache_THEN_get_ports_returns_empty_list(self):
|
||||||
|
"""get_ports should return empty list when no entries exist."""
|
||||||
|
cache = BackendCache()
|
||||||
|
ports = cache.get_ports()
|
||||||
|
self.assertEqual(ports, [])
|
||||||
|
|
||||||
|
def test_GIVEN_cache_with_entries_THEN_get_ports_returns_all_ports(self):
|
||||||
|
"""get_ports should return list of all ports from entries."""
|
||||||
|
cache = BackendCache()
|
||||||
|
|
||||||
|
# Create mock entries with different ports
|
||||||
|
entry1 = Mock()
|
||||||
|
entry1.port = 8000
|
||||||
|
entry2 = Mock()
|
||||||
|
entry2.port = 8001
|
||||||
|
entry3 = Mock()
|
||||||
|
entry3.port = 8002
|
||||||
|
|
||||||
|
cache.entry_list = [entry1, entry2, entry3]
|
||||||
|
|
||||||
|
ports = cache.get_ports()
|
||||||
|
self.assertEqual(ports, [8000, 8001, 8002])
|
||||||
|
|
||||||
|
def test_GIVEN_single_entry_THEN_get_ports_returns_single_port(self):
|
||||||
|
"""get_ports should correctly handle a single entry."""
|
||||||
|
cache = BackendCache()
|
||||||
|
entry = Mock()
|
||||||
|
entry.port = 9000
|
||||||
|
|
||||||
|
cache.entry_list = [entry]
|
||||||
|
|
||||||
|
ports = cache.get_ports()
|
||||||
|
self.assertEqual(ports, [9000])
|
||||||
|
|
||||||
|
|
||||||
|
class TestBackendCacheCheckPath(unittest.TestCase):
|
||||||
|
def setUp(self):
|
||||||
|
self.cache = BackendCache()
|
||||||
|
|
||||||
|
def test_GIVEN_no_matching_path_THEN_check_path_returns_none(self):
|
||||||
|
"""check_path should return None when no entries match the path."""
|
||||||
|
source = Mock()
|
||||||
|
source.name = "test_source"
|
||||||
|
|
||||||
|
cache_entry = Mock()
|
||||||
|
cache_entry.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry.key = Mock()
|
||||||
|
cache_entry.key.source.name = "other_source"
|
||||||
|
cache_entry.key.descriptor = "/some/path"
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_path(source, "/test/path")
|
||||||
|
self.assertIsNone(result)
|
||||||
|
|
||||||
|
def test_GIVEN_terminated_entry_THEN_check_path_ignores_it(self):
|
||||||
|
"""check_path should ignore entries with terminated status."""
|
||||||
|
source = Mock()
|
||||||
|
source.name = "test_source"
|
||||||
|
|
||||||
|
cache_entry = Mock()
|
||||||
|
cache_entry.status = CacheEntryStatus.terminated
|
||||||
|
cache_entry.key = Mock()
|
||||||
|
cache_entry.key.source.name = "test_source"
|
||||||
|
cache_entry.key.descriptor = "/data"
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_path(source, "/data/file.txt")
|
||||||
|
self.assertIsNone(result)
|
||||||
|
|
||||||
|
def test_GIVEN_single_matching_entry_THEN_check_path_returns_it(self):
|
||||||
|
"""check_path should return the matching entry when exactly one matches."""
|
||||||
|
source = Mock()
|
||||||
|
source.name = "test_source"
|
||||||
|
|
||||||
|
cache_entry = Mock()
|
||||||
|
cache_entry.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry.key = Mock()
|
||||||
|
cache_entry.key.source.name = "test_source"
|
||||||
|
cache_entry.key.descriptor = "/data"
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_path(source, "/data/file.txt")
|
||||||
|
self.assertIs(result, cache_entry)
|
||||||
|
|
||||||
|
def test_GIVEN_path_that_does_not_start_with_descriptor_THEN_check_path_returns_none(
|
||||||
|
self,
|
||||||
|
):
|
||||||
|
"""check_path should return None if path doesn't start with descriptor."""
|
||||||
|
source = Mock()
|
||||||
|
source.name = "test_source"
|
||||||
|
|
||||||
|
cache_entry = Mock()
|
||||||
|
cache_entry.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry.key = Mock()
|
||||||
|
cache_entry.key.source.name = "test_source"
|
||||||
|
cache_entry.key.descriptor = "/data"
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_path(source, "/other/file.txt")
|
||||||
|
self.assertIsNone(result)
|
||||||
|
|
||||||
|
def test_GIVEN_multiple_matching_entries_THEN_check_path_raises_exception(self):
|
||||||
|
"""check_path should raise exception when multiple entries match."""
|
||||||
|
source = Mock()
|
||||||
|
source.name = "test_source"
|
||||||
|
|
||||||
|
cache_entry1 = Mock()
|
||||||
|
cache_entry1.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry1.key = Mock()
|
||||||
|
cache_entry1.key.source.name = "test_source"
|
||||||
|
cache_entry1.key.descriptor = "/data"
|
||||||
|
|
||||||
|
cache_entry2 = Mock()
|
||||||
|
cache_entry2.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry2.key = Mock()
|
||||||
|
cache_entry2.key.source.name = "test_source"
|
||||||
|
cache_entry2.key.descriptor = "/data"
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry1, cache_entry2]
|
||||||
|
|
||||||
|
with self.assertRaises(CellxgeneException) as context:
|
||||||
|
self.cache.check_path(source, "/data/file.txt")
|
||||||
|
|
||||||
|
# The CellxgeneException is raised with HTTPStatus as first arg and message as second
|
||||||
|
self.assertEqual(context.exception.message, HTTPStatus.INTERNAL_SERVER_ERROR)
|
||||||
|
self.assertIn("Found 2", context.exception.http_status)
|
||||||
|
|
||||||
|
def test_GIVEN_mixed_entries_THEN_check_path_returns_only_matching_active_entry(
|
||||||
|
self,
|
||||||
|
):
|
||||||
|
"""check_path should correctly filter by source, path, and status."""
|
||||||
|
source = Mock()
|
||||||
|
source.name = "target_source"
|
||||||
|
|
||||||
|
# Terminated entry - should be ignored
|
||||||
|
terminated_entry = Mock()
|
||||||
|
terminated_entry.status = CacheEntryStatus.terminated
|
||||||
|
terminated_entry.key = Mock()
|
||||||
|
terminated_entry.key.source.name = "target_source"
|
||||||
|
terminated_entry.key.descriptor = "/data"
|
||||||
|
|
||||||
|
# Different source - should be ignored
|
||||||
|
other_source_entry = Mock()
|
||||||
|
other_source_entry.status = CacheEntryStatus.loaded
|
||||||
|
other_source_entry.key = Mock()
|
||||||
|
other_source_entry.key.source.name = "other_source"
|
||||||
|
other_source_entry.key.descriptor = "/data"
|
||||||
|
|
||||||
|
# Matching entry - should be returned
|
||||||
|
matching_entry = Mock()
|
||||||
|
matching_entry.status = CacheEntryStatus.loaded
|
||||||
|
matching_entry.key = Mock()
|
||||||
|
matching_entry.key.source.name = "target_source"
|
||||||
|
matching_entry.key.descriptor = "/data"
|
||||||
|
|
||||||
|
self.cache.entry_list = [terminated_entry, other_source_entry, matching_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_path(source, "/data/file.txt")
|
||||||
|
self.assertIs(result, matching_entry)
|
||||||
|
|
||||||
|
|
||||||
|
class TestBackendCacheCheckEntry(unittest.TestCase):
|
||||||
|
def setUp(self):
|
||||||
|
self.cache = BackendCache()
|
||||||
|
|
||||||
|
def test_GIVEN_no_matching_entry_THEN_check_entry_returns_none(self):
|
||||||
|
"""check_entry should return None when no entries match the key."""
|
||||||
|
key = Mock()
|
||||||
|
|
||||||
|
cache_entry = Mock()
|
||||||
|
cache_entry.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry.key = Mock()
|
||||||
|
cache_entry.key.equals.return_value = False
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_entry(key)
|
||||||
|
self.assertIsNone(result)
|
||||||
|
|
||||||
|
def test_GIVEN_terminated_entry_THEN_check_entry_ignores_it(self):
|
||||||
|
"""check_entry should ignore entries with terminated status."""
|
||||||
|
key = Mock()
|
||||||
|
|
||||||
|
cache_entry = Mock()
|
||||||
|
cache_entry.status = CacheEntryStatus.terminated
|
||||||
|
cache_entry.key = Mock()
|
||||||
|
cache_entry.key.equals.return_value = True
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_entry(key)
|
||||||
|
self.assertIsNone(result)
|
||||||
|
|
||||||
|
def test_GIVEN_single_matching_entry_THEN_check_entry_returns_it(self):
|
||||||
|
"""check_entry should return the matching entry when exactly one matches."""
|
||||||
|
key = Mock()
|
||||||
|
|
||||||
|
cache_entry = Mock()
|
||||||
|
cache_entry.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry.key = Mock()
|
||||||
|
cache_entry.key.equals.return_value = True
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_entry(key)
|
||||||
|
self.assertIs(result, cache_entry)
|
||||||
|
|
||||||
|
def test_GIVEN_multiple_matching_entries_THEN_check_entry_raises_exception(self):
|
||||||
|
"""check_entry should raise exception when multiple entries match."""
|
||||||
|
key = Mock()
|
||||||
|
key.dataset = "test_dataset"
|
||||||
|
|
||||||
|
cache_entry1 = Mock()
|
||||||
|
cache_entry1.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry1.key = Mock()
|
||||||
|
cache_entry1.key.equals.return_value = True
|
||||||
|
|
||||||
|
cache_entry2 = Mock()
|
||||||
|
cache_entry2.status = CacheEntryStatus.loaded
|
||||||
|
cache_entry2.key = Mock()
|
||||||
|
cache_entry2.key.equals.return_value = True
|
||||||
|
|
||||||
|
self.cache.entry_list = [cache_entry1, cache_entry2]
|
||||||
|
|
||||||
|
with self.assertRaises(CellxgeneException) as context:
|
||||||
|
self.cache.check_entry(key)
|
||||||
|
|
||||||
|
# The CellxgeneException is raised with HTTPStatus as first arg and message as second
|
||||||
|
self.assertEqual(context.exception.message, HTTPStatus.INTERNAL_SERVER_ERROR)
|
||||||
|
self.assertIn("Found 2", context.exception.http_status)
|
||||||
|
|
||||||
|
def test_GIVEN_mixed_entries_THEN_check_entry_returns_only_matching_active_entry(
|
||||||
|
self,
|
||||||
|
):
|
||||||
|
"""check_entry should correctly filter by key equality and status."""
|
||||||
|
key = Mock()
|
||||||
|
|
||||||
|
# Terminated entry - should be ignored
|
||||||
|
terminated_entry = Mock()
|
||||||
|
terminated_entry.status = CacheEntryStatus.terminated
|
||||||
|
terminated_entry.key = Mock()
|
||||||
|
terminated_entry.key.equals.return_value = True
|
||||||
|
|
||||||
|
# Non-matching entry - should be ignored
|
||||||
|
non_matching_entry = Mock()
|
||||||
|
non_matching_entry.status = CacheEntryStatus.loaded
|
||||||
|
non_matching_entry.key = Mock()
|
||||||
|
non_matching_entry.key.equals.return_value = False
|
||||||
|
|
||||||
|
# Matching entry - should be returned
|
||||||
|
matching_entry = Mock()
|
||||||
|
matching_entry.status = CacheEntryStatus.loaded
|
||||||
|
matching_entry.key = Mock()
|
||||||
|
matching_entry.key.equals.return_value = True
|
||||||
|
|
||||||
|
self.cache.entry_list = [terminated_entry, non_matching_entry, matching_entry]
|
||||||
|
|
||||||
|
result = self.cache.check_entry(key)
|
||||||
|
self.assertIs(result, matching_entry)
|
||||||
|
|
||||||
|
|
||||||
|
class TestBackendCachePrune(unittest.TestCase):
|
||||||
|
def test_GIVEN_entry_in_cache_THEN_prune_removes_it(self):
|
||||||
|
"""prune should remove the entry from the cache."""
|
||||||
|
cache = BackendCache()
|
||||||
|
|
||||||
|
entry_mock = Mock()
|
||||||
|
cache.entry_list = [entry_mock]
|
||||||
|
|
||||||
|
cache.prune(entry_mock)
|
||||||
|
|
||||||
|
self.assertEqual(len(cache.entry_list), 0)
|
||||||
|
self.assertNotIn(entry_mock, cache.entry_list)
|
||||||
|
|
||||||
|
def test_GIVEN_entry_in_cache_THEN_prune_terminates_it(self):
|
||||||
|
"""prune should call terminate on the entry."""
|
||||||
|
cache = BackendCache()
|
||||||
|
|
||||||
|
entry_mock = Mock()
|
||||||
|
cache.entry_list = [entry_mock]
|
||||||
|
|
||||||
|
cache.prune(entry_mock)
|
||||||
|
|
||||||
|
entry_mock.terminate.assert_called_once()
|
||||||
|
|
||||||
|
def test_GIVEN_multiple_entries_THEN_prune_removes_only_target_entry(self):
|
||||||
|
"""prune should only remove the specified entry, not others."""
|
||||||
|
cache = BackendCache()
|
||||||
|
|
||||||
|
entry1 = Mock()
|
||||||
|
entry2 = Mock()
|
||||||
|
entry3 = Mock()
|
||||||
|
|
||||||
|
cache.entry_list = [entry1, entry2, entry3]
|
||||||
|
|
||||||
|
cache.prune(entry2)
|
||||||
|
|
||||||
|
self.assertEqual(len(cache.entry_list), 2)
|
||||||
|
self.assertIn(entry1, cache.entry_list)
|
||||||
|
self.assertNotIn(entry2, cache.entry_list)
|
||||||
|
self.assertIn(entry3, cache.entry_list)
|
||||||
|
|
||||||
|
# Verify only entry2 was terminated
|
||||||
|
entry1.terminate.assert_not_called()
|
||||||
|
entry2.terminate.assert_called_once()
|
||||||
|
entry3.terminate.assert_not_called()
|
||||||
|
|
||||||
|
def test_GIVEN_empty_cache_THEN_prune_raises_value_error(self):
|
||||||
|
"""prune should raise ValueError if entry is not in cache."""
|
||||||
|
cache = BackendCache()
|
||||||
|
|
||||||
|
entry_mock = Mock()
|
||||||
|
|
||||||
|
with self.assertRaises(ValueError):
|
||||||
|
cache.prune(entry_mock)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def connect_mock():
|
||||||
|
with patch("socket.socket") as socketMock:
|
||||||
|
connectMock = socketMock()
|
||||||
|
connectMock.__enter__.return_value = connectMock
|
||||||
|
connectMock.connect_ex.return_value = 1 # Port not in use by default
|
||||||
|
yield connectMock
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def process_backend_mock():
|
||||||
|
with patch("cellxgene_gateway.backend_cache.process_backend") as mock:
|
||||||
|
yield mock
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def sleep_mock():
|
||||||
|
from time import sleep
|
||||||
|
|
||||||
|
with patch("cellxgene_gateway.backend_cache.time.sleep") as mock:
|
||||||
|
mock.side_effect = lambda x: sleep(0.01) # short sleep
|
||||||
|
yield mock
|
||||||
|
|
||||||
|
|
||||||
|
@freeze_time("2025-12-01 10:00:00")
|
||||||
|
def test_GIVEN_empty_cache_THEN_create_entry_uses_starting_port(process_backend_mock):
|
||||||
|
"""create_entry should use port 8000 when cache is empty and port is free."""
|
||||||
|
key = Mock()
|
||||||
|
scripts = ["script1", "script2"]
|
||||||
|
|
||||||
|
cache = BackendCache()
|
||||||
|
result = cache.create_entry(key, scripts)
|
||||||
|
|
||||||
|
# Verify port 8000 was used
|
||||||
|
assert result.port == 8000
|
||||||
|
|
||||||
|
# Verify entry was added to list
|
||||||
|
assert result in cache.entry_list
|
||||||
|
# Verify thread was started
|
||||||
|
process_backend_mock.launch.assert_called_once()
|
||||||
|
|
||||||
|
# Verify CacheEntry was created with correct key
|
||||||
|
assert result.key is key
|
||||||
|
assert result.status == CacheEntryStatus.loading
|
||||||
|
|
||||||
|
|
||||||
|
@freeze_time("2025-12-01 10:00:00")
|
||||||
|
def test_GIVEN_port_in_use_THEN_create_entry_finds_next_available_port(connect_mock):
|
||||||
|
"""create_entry should increment port until finding one that's free."""
|
||||||
|
# Ports 8000 and 8001 are in use, 8002 is free
|
||||||
|
connect_mock.connect_ex.side_effect = [0, 0, 1]
|
||||||
|
|
||||||
|
key = Mock()
|
||||||
|
scripts = []
|
||||||
|
|
||||||
|
cache = BackendCache()
|
||||||
|
result = cache.create_entry(key, scripts)
|
||||||
|
|
||||||
|
# Verify port 8002 was used (after skipping 8000 and 8001)
|
||||||
|
assert result.port == 8002
|
||||||
|
|
||||||
|
|
||||||
|
@freeze_time("2025-12-01 10:00:00")
|
||||||
|
def test_GIVEN_port_exists_in_cache_THEN_create_entry_skips_it():
|
||||||
|
"""create_entry should skip ports already in the cache."""
|
||||||
|
key = Mock()
|
||||||
|
scripts = []
|
||||||
|
|
||||||
|
# Pre-populate cache with port 8000
|
||||||
|
existing_entry = Mock()
|
||||||
|
existing_entry.port = 8000
|
||||||
|
|
||||||
|
cache = BackendCache()
|
||||||
|
cache.entry_list = [existing_entry]
|
||||||
|
|
||||||
|
result = cache.create_entry(key, scripts)
|
||||||
|
|
||||||
|
# Verify port 8001 was used (skipping the existing 8000)
|
||||||
|
assert result.port == 8001
|
||||||
|
|
||||||
|
|
||||||
|
@freeze_time("2025-12-01 10:00:00")
|
||||||
|
def test_GIVEN_create_entry_THEN_background_thread_is_started(process_backend_mock):
|
||||||
|
"""create_entry should start a background thread with correct args."""
|
||||||
|
key = Mock()
|
||||||
|
scripts = ["script1"]
|
||||||
|
|
||||||
|
with patch("cellxgene_gateway.backend_cache.env") as env_mock:
|
||||||
|
env_mock.cellxgene_location = "/path/to/cellxgene"
|
||||||
|
|
||||||
|
cache = BackendCache()
|
||||||
|
result = cache.create_entry(key, scripts)
|
||||||
|
|
||||||
|
# Verify Thread was created with correct target and args
|
||||||
|
process_backend_mock.launch.assert_called_once()
|
||||||
|
call_kwargs = process_backend_mock.launch.call_args[0]
|
||||||
|
assert call_kwargs == ("/path/to/cellxgene", ["script1"], result)
|
||||||
|
|
||||||
|
|
||||||
|
@freeze_time("2025-12-01 10:00:00")
|
||||||
|
def test_GIVEN_create_entry_THEN_entry_is_added_to_cache():
|
||||||
|
"""create_entry should add the created entry to the cache list."""
|
||||||
|
key = Mock()
|
||||||
|
scripts = []
|
||||||
|
|
||||||
|
cache = BackendCache()
|
||||||
|
initial_count = len(cache.entry_list)
|
||||||
|
|
||||||
|
result = cache.create_entry(key, scripts)
|
||||||
|
|
||||||
|
# Verify entry was added
|
||||||
|
assert len(cache.entry_list) == initial_count + 1
|
||||||
|
assert result in cache.entry_list
|
||||||
|
|
||||||
|
|
||||||
|
@freeze_time("2025-12-01 10:00:00")
|
||||||
|
def test_GIVEN_create_entry_THEN_returns_created_entry():
|
||||||
|
"""create_entry should return the created entry."""
|
||||||
|
key = Mock()
|
||||||
|
scripts = []
|
||||||
|
|
||||||
|
cache = BackendCache()
|
||||||
|
result = cache.create_entry(key, scripts)
|
||||||
|
|
||||||
|
# Verify the result is a CacheEntry instance
|
||||||
|
assert isinstance(result, CacheEntry)
|
||||||
|
assert result.status == CacheEntryStatus.loading
|
||||||
|
|||||||
@@ -1,14 +1,16 @@
|
|||||||
import unittest
|
import unittest
|
||||||
|
import tempfile
|
||||||
|
import os
|
||||||
|
import shutil
|
||||||
|
|
||||||
from flask import Flask
|
from flask import Flask
|
||||||
|
|
||||||
from cellxgene_gateway import flask_util
|
from cellxgene_gateway import flask_util
|
||||||
from cellxgene_gateway.cache_entry import CacheEntry, CacheEntryStatus
|
from cellxgene_gateway.cache_entry import CacheEntry, CacheEntryStatus
|
||||||
from cellxgene_gateway.cache_key import CacheKey
|
from cellxgene_gateway.cache_key import CacheKey
|
||||||
from cellxgene_gateway.gateway import app
|
from cellxgene_gateway.items.item import ItemType
|
||||||
from cellxgene_gateway.items.file.fileitem import FileItem
|
from cellxgene_gateway.items.file.fileitem import FileItem
|
||||||
from cellxgene_gateway.items.file.fileitem_source import FileItemSource
|
from cellxgene_gateway.items.file.fileitem_source import FileItemSource
|
||||||
from cellxgene_gateway.items.item import ItemType
|
from cellxgene_gateway.gateway import app
|
||||||
|
|
||||||
key = CacheKey(
|
key = CacheKey(
|
||||||
FileItem("/czi/", name="pbmc3k.h5ad", type=ItemType.h5ad),
|
FileItem("/czi/", name="pbmc3k.h5ad", type=ItemType.h5ad),
|
||||||
@@ -23,6 +25,9 @@ class TestRenderEntry(unittest.TestCase):
|
|||||||
self.app_context.push()
|
self.app_context.push()
|
||||||
self.client = self.app.test_client()
|
self.client = self.app.test_client()
|
||||||
|
|
||||||
|
def tearDown(self):
|
||||||
|
self.app_context.pop()
|
||||||
|
|
||||||
def test_GIVEN_key_and_port_THEN_returns_loading_CacheEntry(self):
|
def test_GIVEN_key_and_port_THEN_returns_loading_CacheEntry(self):
|
||||||
entry = CacheEntry.for_key("some-key", 1)
|
entry = CacheEntry.for_key("some-key", 1)
|
||||||
self.assertEqual(entry.status, CacheEntryStatus.loading)
|
self.assertEqual(entry.status, CacheEntryStatus.loading)
|
||||||
|
|||||||
+24
-2
@@ -1,3 +1,6 @@
|
|||||||
|
import os
|
||||||
|
import shutil
|
||||||
|
import tempfile
|
||||||
import unittest
|
import unittest
|
||||||
from collections import defaultdict
|
from collections import defaultdict
|
||||||
from unittest.mock import patch
|
from unittest.mock import patch
|
||||||
@@ -10,6 +13,7 @@ from cellxgene_gateway.filecrawl import (
|
|||||||
from cellxgene_gateway.items.file.fileitem import FileItem
|
from cellxgene_gateway.items.file.fileitem import FileItem
|
||||||
from cellxgene_gateway.items.file.fileitem_source import FileItemSource
|
from cellxgene_gateway.items.file.fileitem_source import FileItemSource
|
||||||
from cellxgene_gateway.items.item import ItemTree, ItemType
|
from cellxgene_gateway.items.item import ItemTree, ItemType
|
||||||
|
from cellxgene_gateway.gateway import app
|
||||||
|
|
||||||
source = FileItemSource("/tmp")
|
source = FileItemSource("/tmp")
|
||||||
|
|
||||||
@@ -25,6 +29,14 @@ def make_entry(subpath="somepath", annotations=None):
|
|||||||
|
|
||||||
|
|
||||||
class TestRenderEntry(unittest.TestCase):
|
class TestRenderEntry(unittest.TestCase):
|
||||||
|
def setUp(self):
|
||||||
|
self.app = app
|
||||||
|
self.app_context = self.app.test_request_context()
|
||||||
|
self.app_context.push()
|
||||||
|
|
||||||
|
def tearDown(self):
|
||||||
|
self.app_context.pop()
|
||||||
|
|
||||||
def test_GIVEN_path_both_slash_THEN_view_has_single_slash(self):
|
def test_GIVEN_path_both_slash_THEN_view_has_single_slash(self):
|
||||||
entry = make_entry(subpath="/somepath/")
|
entry = make_entry(subpath="/somepath/")
|
||||||
rendered = render_item(entry, source)
|
rendered = render_item(entry, source)
|
||||||
@@ -47,6 +59,15 @@ class TestRenderEntry(unittest.TestCase):
|
|||||||
|
|
||||||
|
|
||||||
class TestRenderAnnotation(unittest.TestCase):
|
class TestRenderAnnotation(unittest.TestCase):
|
||||||
|
|
||||||
|
def setUp(self):
|
||||||
|
self.app = app
|
||||||
|
self.app_context = self.app.test_request_context()
|
||||||
|
self.app_context.push()
|
||||||
|
|
||||||
|
def tearDown(self):
|
||||||
|
self.app_context.pop()
|
||||||
|
|
||||||
@patch("cellxgene_gateway.filecrawl.enable_annotations", new=True)
|
@patch("cellxgene_gateway.filecrawl.enable_annotations", new=True)
|
||||||
def test_GIVEN_no_annotation_THEN_new_alone(self):
|
def test_GIVEN_no_annotation_THEN_new_alone(self):
|
||||||
entry = make_entry(annotations=None)
|
entry = make_entry(annotations=None)
|
||||||
@@ -103,12 +124,13 @@ class TestRenderItemSource(unittest.TestCase):
|
|||||||
|
|
||||||
class TestRenderItemTree(unittest.TestCase):
|
class TestRenderItemTree(unittest.TestCase):
|
||||||
def setUp(self):
|
def setUp(self):
|
||||||
from cellxgene_gateway.gateway import app
|
|
||||||
|
|
||||||
self.app = app
|
self.app = app
|
||||||
self.app_context = self.app.test_request_context()
|
self.app_context = self.app.test_request_context()
|
||||||
self.app_context.push()
|
self.app_context.push()
|
||||||
|
|
||||||
|
def tearDown(self):
|
||||||
|
self.app_context.pop()
|
||||||
|
|
||||||
@patch("cellxgene_gateway.items.file.fileitem_source.FileItemSource")
|
@patch("cellxgene_gateway.items.file.fileitem_source.FileItemSource")
|
||||||
def test_GIVEN_deep_nested_dirs_THEN_includes_dirs_in_output(self, item_source):
|
def test_GIVEN_deep_nested_dirs_THEN_includes_dirs_in_output(self, item_source):
|
||||||
item_source.name = "FakeSource"
|
item_source.name = "FakeSource"
|
||||||
|
|||||||
@@ -0,0 +1,54 @@
|
|||||||
|
import json
|
||||||
|
import unittest
|
||||||
|
from types import SimpleNamespace
|
||||||
|
|
||||||
|
from cellxgene_gateway.gateway import do_GET_status_json, app, cache
|
||||||
|
from cellxgene_gateway.cache_entry import CacheEntry, CacheEntryStatus
|
||||||
|
|
||||||
|
|
||||||
|
class TestGatewayStatusJson(unittest.TestCase):
|
||||||
|
def test_do_GET_status_json_returns_expected_structure(self):
|
||||||
|
# Create a minimal fake key with required attributes
|
||||||
|
h5ad_item = SimpleNamespace(descriptor="somedir/dataset.h5ad")
|
||||||
|
key = SimpleNamespace(
|
||||||
|
h5ad_item=h5ad_item,
|
||||||
|
annotation_descriptor="somedir/dataset_annotations/foo.csv",
|
||||||
|
)
|
||||||
|
|
||||||
|
# Create a CacheEntry with known launchtime/timestamp/status
|
||||||
|
entry = CacheEntry(
|
||||||
|
None,
|
||||||
|
key,
|
||||||
|
8000,
|
||||||
|
111,
|
||||||
|
222,
|
||||||
|
CacheEntryStatus.loaded,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
|
||||||
|
# Install into the gateway cache and set app launchtime
|
||||||
|
cache.entry_list = [entry]
|
||||||
|
app.extensions.setdefault("cellxgene_gateway", {})["launchtime"] = "LAUNCH_TIME"
|
||||||
|
|
||||||
|
rv = do_GET_status_json()
|
||||||
|
|
||||||
|
data = json.loads(rv)
|
||||||
|
# top-level launchtime comes from app.extensions
|
||||||
|
self.assertEqual("LAUNCH_TIME", data["launchtime"])
|
||||||
|
|
||||||
|
self.assertIn("entry_list", data)
|
||||||
|
self.assertEqual(1, len(data["entry_list"]))
|
||||||
|
|
||||||
|
e = data["entry_list"][0]
|
||||||
|
self.assertEqual("somedir/dataset.h5ad", e["dataset"])
|
||||||
|
self.assertEqual("somedir/dataset_annotations/foo.csv", e["annotation_file"])
|
||||||
|
self.assertEqual("loaded", e["status"])
|
||||||
|
self.assertEqual(111, e["launchtime"])
|
||||||
|
self.assertEqual(222, e["last_access"])
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
Reference in New Issue
Block a user