Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions config/config.default.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,8 @@ mongo:
authsource:
attempts: 3 # Number of attempts to accomplish operation before considering it failed.
shutdown_timeout: 3 # Time in seconds to wait for serialize task to finish.
compression_level: 10 # zstd level used to compress stored job documents. Higher means greater compression, at greater CPU cost.
max_stored_doc_bytes: 1073741824 # Compressed job-docs at or above this size are not persisted.
tier0:
backend: gandalf
backend_infores: infores:dogpark-tier0
Expand Down
12 changes: 12 additions & 0 deletions src/retriever/config/general.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,18 @@ class MongoSettings(BaseModel):
shutdown_timeout: Annotated[
int, Field(description="Time in seconds to wait for serialize task to finish.")
] = 3
compression_level: Annotated[
int,
Field(
description="zstd level used to compress stored job documents. Higher means greater compression, at greater CPU cost."
),
] = 10
max_stored_doc_bytes: Annotated[
int,
Field(
description="Compressed job-docs at or above this size are not persisted."
),
] = 1024**3


class TelemetrySettings(BaseModel):
Expand Down
35 changes: 20 additions & 15 deletions src/retriever/lookup/lookup.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@
MONGO_QUEUE = MongoQueue()
OP_TABLE_MANAGER = OpTableManager()

ZSTD_COMPRESSOR = zstandard.ZstdCompressor()
ZSTD_COMPRESSOR = zstandard.ZstdCompressor(level=CONFIG.mongo.compression_level)

CALLBACK_BODY_LOG_LIMIT = 500
"""Truncate logged callback response bodies to this many chars."""
Expand Down Expand Up @@ -172,6 +172,10 @@ async def lookup(query: QueryInfo) -> tuple[HTTPStatus, ResponseDict]:
job_log.info(
f"Begin processing job {job_id} for client {get_submitter(query)}."
)
if (query.body.get("parameters") or {}).get("dehydrated"):
job_log.info(
"Query running in dehydrated mode. Response will not be stored."
)

qgraph = query.body["message"].get("query_graph")
if qgraph is None:
Expand Down Expand Up @@ -398,21 +402,22 @@ def tracked_response(

# Cast because TypedDict has some *really annoying* interactions with more general dicts
kgraph = response["message"].get("knowledge_graph") or {}
# Don't store dehydrated job responses, they're usually huge
dehydrated = bool(((query.body or {}).get("parameters") or {}).get("dehydrated"))
state = ResponseState(
job_id=query.job_id,
event_time=datetime.now().astimezone(),
knodes=len(kgraph.get("nodes") or {}),
kedges=len(kgraph.get("edges") or {}),
aux_graphs=len(response["message"].get("auxiliary_graphs") or {}),
results=len(response["message"].get("results") or []),
status=(response.get("status") or "Running"),
description=response.get("description"),
)
if not dehydrated:
state["response"] = ZSTD_COMPRESSOR.compress(ormsgpack.packb(response))
try:
MONGO_QUEUE.put(
"job_state",
ResponseState(
job_id=query.job_id,
response=ZSTD_COMPRESSOR.compress(ormsgpack.packb(response)),
event_time=datetime.now().astimezone(),
knodes=len(kgraph.get("nodes") or {}),
kedges=len(kgraph.get("edges") or {}),
aux_graphs=len(response["message"].get("auxiliary_graphs") or {}),
results=len(response["message"].get("results") or []),
status=(response.get("status") or "Running"),
description=response.get("description"),
),
)
MONGO_QUEUE.put("job_state", state)
except MongoOutage:
response["logs"] = [
*(response.get("logs") or []),
Expand Down
5 changes: 4 additions & 1 deletion src/retriever/query.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@

tracer = trace.get_tracer("lookup.execution.tracer")
MONGO_QUEUE = MongoQueue()
ZSTD_COMPRESSOR = zstandard.ZstdCompressor()
ZSTD_COMPRESSOR = zstandard.ZstdCompressor(level=CONFIG.mongo.compression_level)
ZSTD_DECOMPRESSOR = zstandard.ZstdDecompressor()


Expand Down Expand Up @@ -119,6 +119,9 @@ def _record_initial_state(
is_async=ctx.background_tasks is not None,
worker_pid=worker.get_pid(),
worker_started_at=worker.get_started_at(),
dehydrated=bool(
((query.body or {}).get("parameters") or {}).get("dehydrated")
),
**{
k: v
for k, v in query_metadata._asdict().items()
Expand Down
Loading
Loading