Skip to content

Commit f623888

Browse files
committed
feat(storage): implement support for storing artifacts to S3 (closes
1 parent 224428c commit f623888

14 files changed

Lines changed: 414 additions & 38 deletions

File tree

README.md

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ A simplistic, opinionated remote update server implementing hawkBit™'s [DDI AP
2525
```
2626

2727
3. Launch gooseBit:
28+
2829
```bash
2930
python main.py
3031
```
@@ -140,16 +141,17 @@ The structure of gooseBit is as follows:
140141

141142
- `api`: Files for the API.
142143
- `ui`: Files for the UI.
143-
- `bff`: Backend for frontend API.
144-
- `static`: Static files.
145-
- `templates`: Jinja2 formatted templates.
146-
- `nav`: Navbar handler.
144+
- `bff`: Backend for frontend API.
145+
- `static`: Static files.
146+
- `templates`: Jinja2 formatted templates.
147+
- `nav`: Navbar handler.
147148
- `updater`: DDI API handler and device update manager.
148149
- `updates`: SWUpdate file parsing.
149150
- `auth`: Authentication functions and permission handling.
150151
- `models`: Database models.
151152
- `db`: Database config and initialization.
152153
- `schema`: Pydantic models used for API type hinting.
153154
- `settings`: Settings loader and handler.
155+
- `storage`: Storage for software artifacts.
154156
- `telemetry`: Telemetry data handlers.
155157
- `routes`: Routes for a giving endpoint, including the router.

conftest.py

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,14 +36,22 @@ async def clear_cache():
3636
@pytest_asyncio.fixture(scope="function")
3737
async def test_app():
3838
from goosebit.users import create_initial_user
39+
from goosebit.storage import init_storage, close_storage
40+
from goosebit.settings import config
3941

4042
async with RegisterTortoise(
4143
app=app,
4244
config=TORTOISE_CONF,
4345
):
4446
await Tortoise.generate_schemas()
4547
await create_initial_user(username="testing@goosebit.test", hashed_pwd=PWD_CXT.hash("test"))
46-
yield app
48+
49+
init_storage(config)
50+
51+
try:
52+
yield app
53+
finally:
54+
close_storage()
4755

4856

4957
@pytest_asyncio.fixture(scope="function")

goosebit.yaml

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,3 +68,13 @@ logging:
6868
root:
6969
level: INFO
7070
handlers: [console]
71+
72+
# Storage settings (default backend is filesystem)
73+
# storage:
74+
# backend: s3
75+
# s3:
76+
# bucket: goosebit
77+
# # region: us-east-1
78+
# # endpoint_url: http://localhost:9000 # example for self-hosted min.io
79+
# access_key_id: minioadmin
80+
# secret_access_key: minioadmin

goosebit/__init__.py

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
from goosebit import api, db, ui, updater
1717
from goosebit.auth import get_user_from_request, login_user, redirect_if_authenticated
1818
from goosebit.settings import PWD_CXT, config
19+
from goosebit.storage import init_storage, close_storage
1920
from goosebit.ui.nav import nav
2021
from goosebit.ui.static import static
2122
from goosebit.ui.templates import templates
@@ -29,9 +30,15 @@ async def lifespan(_: FastAPI):
2930
db_ready = await db.init()
3031
if not db_ready:
3132
logger.exception("DB does not exist, try running `poetry run aerich upgrade`.")
33+
34+
init_storage(config)
35+
logger.debug(f"Initialized storage backend: {config.storage.backend}")
36+
3237
if db_ready:
3338
yield
39+
3440
await db.close()
41+
close_storage()
3542

3643

3744
app = FastAPI(

goosebit/api/v1/download/routes.py

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,9 @@
11
from fastapi import APIRouter, HTTPException
22
from fastapi.requests import Request
3-
from fastapi.responses import FileResponse, RedirectResponse
3+
from fastapi.responses import FileResponse, RedirectResponse, StreamingResponse
44

55
from goosebit.db.models import Software
6+
from goosebit.storage import get_storage
67

78
router = APIRouter(prefix="/download", tags=["download"])
89

@@ -18,5 +19,16 @@ async def download_file(_: Request, file_id: int):
1819
media_type="application/octet-stream",
1920
filename=software.path.name,
2021
)
21-
else:
22-
return RedirectResponse(url=software.uri)
22+
23+
storage = get_storage()
24+
try:
25+
url = await storage.get_download_url(software.uri)
26+
return RedirectResponse(url=url)
27+
except Exception:
28+
# Fallback to streaming if redirect fails.
29+
file_stream = storage.get_file_stream(software.uri)
30+
return StreamingResponse(
31+
file_stream,
32+
media_type="application/octet-stream",
33+
headers={"Content-Disposition": f"attachment; filename={software.path.name}"}
34+
)

goosebit/settings/schema.py

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,21 @@ class PrometheusSettings(BaseModel):
3434
class MetricsSettings(BaseModel):
3535
prometheus: PrometheusSettings = PrometheusSettings()
3636

37+
class StorageType(StrEnum):
38+
FILESYSTEM = "filesystem"
39+
S3 = "s3"
40+
41+
class S3StorageSettings(BaseModel):
42+
bucket: str
43+
region: str = "us-east-1"
44+
endpoint_url: str | None = None
45+
access_key_id: str | None = None
46+
secret_access_key: str | None = None
47+
48+
49+
class StorageSettings(BaseModel):
50+
backend: StorageType = StorageType.FILESYSTEM
51+
s3: S3StorageSettings | None = None
3752

3853
class GooseBitSettings(BaseSettings):
3954
model_config = SettingsConfigDict(env_prefix="GOOSEBIT_")
@@ -54,6 +69,8 @@ class GooseBitSettings(BaseSettings):
5469

5570
artifacts_dir: Path = GOOSEBIT_ROOT_DIR.joinpath("artifacts")
5671

72+
storage: StorageSettings = StorageSettings()
73+
5774
metrics: MetricsSettings = MetricsSettings()
5875

5976
logging: dict = LOGGING_DEFAULT

goosebit/storage/__init__.py

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
1+
from abc import ABC, abstractmethod
2+
from pathlib import Path
3+
from typing import AsyncGenerator
4+
5+
from goosebit.settings.schema import GooseBitSettings, StorageType
6+
7+
8+
class Storage(ABC):
9+
@abstractmethod
10+
async def store_file(self, source_path: Path, key: str) -> str:
11+
pass
12+
13+
@abstractmethod
14+
async def get_file_stream(self, uri: str) -> AsyncGenerator[bytes, None]:
15+
pass
16+
17+
@abstractmethod
18+
async def get_download_url(self, uri: str) -> str:
19+
pass
20+
21+
22+
_storage: Storage | None = None
23+
24+
25+
def init_storage(config: GooseBitSettings) -> None:
26+
global _storage
27+
_storage = create_storage(config)
28+
29+
def create_storage(config: GooseBitSettings) -> Storage:
30+
from .filesystem import FilesystemStorage
31+
from .s3 import S3Storage
32+
33+
if config.storage.backend == StorageType.FILESYSTEM:
34+
return FilesystemStorage(base_path=config.artifacts_dir)
35+
36+
elif config.storage.backend == StorageType.S3:
37+
if config.storage.s3 is None:
38+
return FilesystemStorage(base_path=config.artifacts_dir)
39+
40+
config = config.storage.s3
41+
return S3Storage(
42+
bucket=config.bucket,
43+
region=config.region,
44+
endpoint_url=config.endpoint_url,
45+
access_key_id=config.access_key_id,
46+
secret_access_key=config.secret_access_key
47+
)
48+
49+
else:
50+
raise ValueError(f"Unknown storage backend type: {config.storage.backend}")
51+
52+
def get_storage() -> Storage:
53+
if _storage is None:
54+
raise RuntimeError("Storage backend not initialized. Call init_storage() first.")
55+
56+
return _storage
57+
58+
def close_storage() -> None:
59+
global _storage
60+
_storage = None

goosebit/storage/filesystem.py

Lines changed: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,71 @@
1+
import shutil
2+
from pathlib import Path
3+
from typing import AsyncGenerator
4+
from urllib.parse import urlparse
5+
6+
import httpx
7+
from anyio import open_file
8+
9+
from . import Storage
10+
11+
12+
class FilesystemStorage(Storage):
13+
def __init__(self, base_path: Path):
14+
self.base_path = Path(base_path)
15+
self.base_path.mkdir(parents=True, exist_ok=True)
16+
17+
async def store_file(self, source_path: Path, key: str) -> str:
18+
dest_path = self.base_path / key
19+
dest_path.parent.mkdir(parents=True, exist_ok=True)
20+
21+
shutil.copy2(source_path, dest_path)
22+
23+
return dest_path.resolve().as_uri()
24+
25+
async def get_file_stream(self, uri: str) -> AsyncGenerator[bytes, None]:
26+
parsed = urlparse(uri)
27+
28+
if parsed.scheme in ("http", "https"):
29+
async with httpx.AsyncClient() as client:
30+
async with client.stream("GET", uri) as response:
31+
response.raise_for_status()
32+
async for chunk in response.aiter_bytes(8192):
33+
yield chunk
34+
35+
elif parsed.scheme == "file":
36+
file_path = self._extract_path_from_uri(uri)
37+
if not file_path.exists():
38+
raise FileNotFoundError(f"File not found: {file_path}")
39+
40+
async with await open_file(file_path, "rb") as f:
41+
while True:
42+
chunk = await f.read(8192)
43+
if not chunk:
44+
break
45+
yield chunk
46+
else:
47+
raise ValueError(f"Unsupported URI scheme '{parsed.scheme}' for filesystem backend: {uri}")
48+
49+
async def get_download_url(self, uri: str) -> str:
50+
parsed = urlparse(uri)
51+
52+
if parsed.scheme in ("http", "https"):
53+
return uri
54+
55+
elif parsed.scheme == "file":
56+
file_path = self._extract_path_from_uri(uri)
57+
if not file_path.exists():
58+
raise FileNotFoundError(f"File not found: {file_path}")
59+
60+
return file_path.resolve().as_uri()
61+
62+
else:
63+
raise ValueError(f"Unsupported URI scheme '{parsed.scheme}' for filesystem backend: {uri}")
64+
65+
def _extract_path_from_uri(self, uri: str) -> Path:
66+
parsed = urlparse(uri)
67+
68+
if parsed.scheme != "file":
69+
raise ValueError(f"Expected file:// URI, got: {uri}")
70+
71+
return Path(parsed.path)

goosebit/storage/s3.py

Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,98 @@
1+
import asyncio
2+
from pathlib import Path
3+
from typing import AsyncGenerator
4+
5+
from boto3.session import Session
6+
from botocore.config import Config
7+
from botocore.exceptions import ClientError
8+
9+
from . import Storage
10+
11+
class S3Storage(Storage):
12+
def __init__(
13+
self,
14+
bucket: str,
15+
region: str = "us-east-1",
16+
endpoint_url: str | None = None,
17+
access_key_id: str | None = None,
18+
secret_access_key: str | None = None
19+
):
20+
self.bucket = bucket
21+
22+
config = Config(
23+
region_name = region,
24+
connect_timeout = 10,
25+
read_timeout = 60,
26+
retries = {"max_attempts": 5, "mode": "adaptive"},
27+
signature_version = "s3v4",
28+
)
29+
30+
session_config = {}
31+
if access_key_id is not None:
32+
session_config["aws_access_key_id"] = access_key_id
33+
if secret_access_key is not None:
34+
session_config["aws_secret_access_key"] = secret_access_key
35+
36+
session = Session(**session_config)
37+
38+
self.s3_client = session.client("s3", config=config, endpoint_url=endpoint_url)
39+
40+
async def store_file(self, source_path: Path, key: str) -> str:
41+
try:
42+
loop = asyncio.get_event_loop()
43+
await loop.run_in_executor(
44+
None,
45+
self.s3_client.upload_file,
46+
str(source_path),
47+
self.bucket,
48+
key
49+
)
50+
return f"s3://{self.bucket}/{key}"
51+
except ClientError as e:
52+
raise ValueError(f"S3 upload failed: {e}")
53+
54+
async def get_file_stream(self, uri: str) -> AsyncGenerator[bytes, None]:
55+
key = self._extract_key_from_uri(uri)
56+
57+
try:
58+
loop = asyncio.get_event_loop()
59+
response = await loop.run_in_executor(
60+
None,
61+
self.s3_client.get_object,
62+
self.bucket,
63+
key
64+
)
65+
66+
body = response["Body"]
67+
try:
68+
while True:
69+
chunk = await loop.run_in_executor(None, body.read, 8192)
70+
if not chunk:
71+
break
72+
yield chunk
73+
finally:
74+
await loop.run_in_executor(None, body.close)
75+
76+
except ClientError as e:
77+
raise ValueError(f"S3 download failed: {e}")
78+
79+
async def get_download_url(self, uri: str) -> str:
80+
key = self._extract_key_from_uri(uri)
81+
82+
try:
83+
loop = asyncio.get_event_loop()
84+
return await loop.run_in_executor(
85+
None,
86+
self.s3_client.generate_presigned_url,
87+
"get_object",
88+
{"Bucket": self.bucket, "Key": key},
89+
3600 # 1 hour expiration
90+
)
91+
except ClientError as e:
92+
raise ValueError(f"Failed to generate presigned URL: {e}")
93+
94+
def _extract_key_from_uri(self, uri: str) -> str:
95+
if not uri.startswith(f"s3://{self.bucket}/"):
96+
raise ValueError(f"Invalid S3 URI for bucket {self.bucket}: {uri}")
97+
98+
return uri.replace(f"s3://{self.bucket}/", "")

0 commit comments

Comments
 (0)