Skip to content

Commit 43018ac

Browse files
committed
fixing lint
1 parent da4e6ab commit 43018ac

File tree

3 files changed

+19
-12
lines changed

3 files changed

+19
-12
lines changed

backend/app/dependencies.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -93,7 +93,9 @@ async def get_rabbitmq() -> AbstractChannel:
9393
def get_blocking_rabbitmq() -> BlockingChannel:
9494
"""Legacy blocking RabbitMQ client (for extractors that need it)"""
9595
credentials = pika.PlainCredentials(settings.RABBITMQ_USER, settings.RABBITMQ_PASS)
96-
parameters = pika.ConnectionParameters(settings.RABBITMQ_HOST, credentials=credentials)
96+
parameters = pika.ConnectionParameters(
97+
settings.RABBITMQ_HOST, credentials=credentials
98+
)
9799
connection = pika.BlockingConnection(parameters)
98100
return connection.channel()
99101

backend/app/rabbitmq/listeners.py

Lines changed: 11 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -22,11 +22,15 @@
2222

2323

2424
async def create_reply_queue(channel: AbstractChannel):
25-
if (config_entry := await ConfigEntryDB.find_one({"key": "instance_id"})) is not None:
25+
if (
26+
config_entry := await ConfigEntryDB.find_one({"key": "instance_id"})
27+
) is not None:
2628
instance_id = config_entry.value
2729
else:
2830
instance_id = "".join(
29-
random.choice(string.ascii_uppercase + string.ascii_lowercase + string.digits)
31+
random.choice(
32+
string.ascii_uppercase + string.ascii_lowercase + string.digits
33+
)
3034
for _ in range(10)
3135
)
3236
config_entry = ConfigEntryDB(key="instance_id", value=instance_id)
@@ -36,7 +40,9 @@ async def create_reply_queue(channel: AbstractChannel):
3640

3741
# Use aio_pika methods instead of pika methods
3842
exchange = await channel.declare_exchange("clowder", durable=True)
39-
queue = await channel.declare_queue(queue_name, durable=True, exclusive=False, auto_delete=False)
43+
queue = await channel.declare_queue(
44+
queue_name, durable=True, exclusive=False, auto_delete=False
45+
)
4046
await queue.bind(exchange)
4147

4248
return queue.name
@@ -61,7 +67,6 @@ async def submit_file_job(
6167
)
6268
await job.insert()
6369

64-
6570
current_secretKey = await get_user_job_key(user.email)
6671
msg_body = EventListenerJobMessage(
6772
filename=file_out.name,
@@ -79,7 +84,7 @@ async def submit_file_job(
7984
print("RABBITMQ_CLIENT: " + str(rabbitmq_client))
8085
await rabbitmq_client.default_exchange.publish(
8186
aio_pika.Message(
82-
body=json.dumps(msg_body.dict(), ensure_ascii=False).encode('utf-8'),
87+
body=json.dumps(msg_body.dict(), ensure_ascii=False).encode("utf-8"),
8388
content_type="application/json",
8489
delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
8590
reply_to=reply_to,
@@ -117,7 +122,7 @@ async def submit_dataset_job(
117122
reply_to = await create_reply_queue(rabbitmq_client)
118123
await rabbitmq_client.default_exchange.publish(
119124
aio_pika.Message(
120-
body=json.dumps(msg_body.dict(), ensure_ascii=False).encode('utf-8'),
125+
body=json.dumps(msg_body.dict(), ensure_ascii=False).encode("utf-8"),
121126
content_type="application/json",
122127
delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
123128
reply_to=reply_to,

backend/app/routers/files.py

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@
4141
router = APIRouter()
4242
security = HTTPBearer()
4343

44+
4445
class CustomJSONEncoder(JSONEncoder):
4546
def default(self, obj):
4647
if isinstance(obj, PydanticObjectId):
@@ -148,13 +149,12 @@ async def add_file_entry(
148149

149150
# Publish a message when indexing is complete
150151

151-
152152
# FIXED: Use aio_pika publishing
153153
message_body = {
154154
"event_type": "file_indexed",
155155
"file_data": json.loads(new_file.json()),
156156
"user": json.loads(user.json()),
157-
"timestamp": datetime.now().isoformat()
157+
"timestamp": datetime.now().isoformat(),
158158
}
159159

160160
# Get the exchange first
@@ -163,7 +163,7 @@ async def add_file_entry(
163163
# Use aio_pika publish method
164164
await exchange.publish(
165165
aio_pika.Message(
166-
body=json.dumps(message_body).encode('utf-8'),
166+
body=json.dumps(message_body).encode("utf-8"),
167167
content_type="application/json",
168168
delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
169169
),
@@ -201,7 +201,7 @@ async def add_local_file_entry(
201201
"event_type": "file_indexed",
202202
"file_data": json.loads(new_file.json()),
203203
"user": json.loads(user.json()),
204-
"timestamp": datetime.now().isoformat()
204+
"timestamp": datetime.now().isoformat(),
205205
}
206206

207207
# Get the exchange first
@@ -210,7 +210,7 @@ async def add_local_file_entry(
210210
# Use aio_pika publish method
211211
await exchange.publish(
212212
aio_pika.Message(
213-
body=json.dumps(message_body).encode('utf-8'),
213+
body=json.dumps(message_body).encode("utf-8"),
214214
content_type="application/json",
215215
delivery_mode=aio_pika.DeliveryMode.PERSISTENT,
216216
),

0 commit comments

Comments
 (0)