Publish #1
@@ -5,8 +5,10 @@ description = ""
|
|||||||
readme = "README.md"
|
readme = "README.md"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
|
"aiobotocore>=2.26.0",
|
||||||
"dotenv>=0.9.9",
|
"dotenv>=0.9.9",
|
||||||
"fastapi>=0.121.1",
|
"fastapi>=0.121.1",
|
||||||
|
"minio>=7.2.20",
|
||||||
"python-multipart>=0.0.20",
|
"python-multipart>=0.0.20",
|
||||||
"uvicorn>=0.38.0",
|
"uvicorn>=0.38.0",
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -0,0 +1,88 @@
|
|||||||
|
import os
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from dotenv import load_dotenv
|
||||||
|
from aiobotocore.session import get_session
|
||||||
|
|
||||||
|
|
||||||
|
class S3Worker:
|
||||||
|
def __init__(self):
|
||||||
|
load_dotenv()
|
||||||
|
|
||||||
|
self._S3_ACCESS_KEY_ID = os.getenv('S3_ACCESS_KEY_ID')
|
||||||
|
self._S3_SECRET_ACCESS_KEY = os.getenv('S3_SECRET_ACCESS_KEY')
|
||||||
|
self._S3_ENDPOINT_URL = os.getenv('S3_ENDPOINT_URL')
|
||||||
|
|
||||||
|
self._s3_session = get_session()
|
||||||
|
self._s3_client = None
|
||||||
|
|
||||||
|
self.bucket = os.getenv('BUCKET_NAME')
|
||||||
|
|
||||||
|
async def __aenter__(self):
|
||||||
|
self._client = await self._s3_session.create_client(
|
||||||
|
"s3",
|
||||||
|
region_name="us-east-1",
|
||||||
|
aws_access_key_id=self._S3_ACCESS_KEY_ID,
|
||||||
|
aws_secret_access_key=self._S3_SECRET_ACCESS_KEY,
|
||||||
|
endpoint_url=self._S3_ENDPOINT_URL,
|
||||||
|
).__aenter__()
|
||||||
|
return self
|
||||||
|
|
||||||
|
async def __aexit__(self, exc_type, exc, tb):
|
||||||
|
await self._client.__aexit__(exc_type, exc, tb)
|
||||||
|
|
||||||
|
|
||||||
|
async def upload_file(self, key: str, data: bytes):
|
||||||
|
await self._client.put_object(
|
||||||
|
Bucket=self.bucket,
|
||||||
|
Key=key,
|
||||||
|
Body=data,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def generate_upload_url(self, key: str, content_type: str, expires_in: int = 300) -> str:
|
||||||
|
return await self._client.generate_presigned_url("put_object",
|
||||||
|
Params={
|
||||||
|
"Bucket": self.bucket,
|
||||||
|
"Key": key,
|
||||||
|
"ContentType": content_type,
|
||||||
|
},
|
||||||
|
ExpiresIn=expires_in,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def download_file(self, key: str):
|
||||||
|
return await self._client.get_object(
|
||||||
|
Bucket=self.bucket,
|
||||||
|
Key=key,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def generate_download_url(self, key: str, filename: str, expires_in: int = 300) -> str:
|
||||||
|
return await self._client.generate_presigned_url("get_object",
|
||||||
|
Params={
|
||||||
|
"Bucket": self.bucket,
|
||||||
|
"Key": key,
|
||||||
|
"ResponseContentDisposition": (
|
||||||
|
f'attachment; filename="{filename}"'
|
||||||
|
),
|
||||||
|
},
|
||||||
|
ExpiresIn=expires_in,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def test_run():
|
||||||
|
async with S3Worker() as worker:
|
||||||
|
await worker.upload_file("test.txt", b"hello")
|
||||||
|
file = await worker.download_file("test.txt")
|
||||||
|
file_text = await file["Body"].read()
|
||||||
|
print(file_text)
|
||||||
|
|
||||||
|
url = await worker.generate_download_url("test.txt", "hello.txt")
|
||||||
|
print(url)
|
||||||
|
url = await worker.generate_upload_url("some.jpg", content_type="image/jpeg")
|
||||||
|
print(url)
|
||||||
|
url = await worker.generate_download_url("some.jpg", "image.jpg")
|
||||||
|
print(url)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
asyncio.run(test_run())
|
||||||
Reference in New Issue
Block a user