mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-27 06:01:01 +02:00
fix: speed up config startup and prevent timeout cascades (#1062)
- Increase per-request config timeout from 10s to 60s - Parallelize per-workspace config fetches with asyncio.gather - Track applied workspaces so retries resume instead of restarting With 50 workspaces, sequential fetches with 10s timeouts caused a cascade: timeouts triggered retries that re-fetched everything, resulting in 10+ minute startup times.
This commit is contained in:
parent
33129eaf99
commit
35802442fb
1 changed files with 29 additions and 7 deletions
|
|
@ -139,7 +139,7 @@ class AsyncProcessor:
|
||||||
workspace=workspace,
|
workspace=workspace,
|
||||||
type=config_type,
|
type=config_type,
|
||||||
),
|
),
|
||||||
timeout=10,
|
timeout=60,
|
||||||
)
|
)
|
||||||
if resp.error:
|
if resp.error:
|
||||||
raise RuntimeError(f"Config error: {resp.error.message}")
|
raise RuntimeError(f"Config error: {resp.error.message}")
|
||||||
|
|
@ -156,7 +156,7 @@ class AsyncProcessor:
|
||||||
operation="getkeys-all-ws",
|
operation="getkeys-all-ws",
|
||||||
type=config_type,
|
type=config_type,
|
||||||
),
|
),
|
||||||
timeout=10,
|
timeout=60,
|
||||||
)
|
)
|
||||||
if resp.error:
|
if resp.error:
|
||||||
raise RuntimeError(f"Config error: {resp.error.message}")
|
raise RuntimeError(f"Config error: {resp.error.message}")
|
||||||
|
|
@ -168,10 +168,21 @@ class AsyncProcessor:
|
||||||
for v in resp.values:
|
for v in resp.values:
|
||||||
workspaces.add(v.workspace)
|
workspaces.add(v.workspace)
|
||||||
|
|
||||||
# Step 2: fetch values per workspace
|
# Step 2: fetch values per workspace in parallel
|
||||||
grouped = {}
|
async def fetch_one(ws):
|
||||||
for ws in workspaces:
|
|
||||||
kv = await self._fetch_type_workspace(client, ws, config_type)
|
kv = await self._fetch_type_workspace(client, ws, config_type)
|
||||||
|
return ws, kv
|
||||||
|
|
||||||
|
results = await asyncio.gather(
|
||||||
|
*(fetch_one(ws) for ws in workspaces),
|
||||||
|
return_exceptions=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
grouped = {}
|
||||||
|
for result in results:
|
||||||
|
if isinstance(result, Exception):
|
||||||
|
raise result
|
||||||
|
ws, kv = result
|
||||||
if kv:
|
if kv:
|
||||||
grouped[ws] = kv
|
grouped[ws] = kv
|
||||||
|
|
||||||
|
|
@ -194,7 +205,14 @@ class AsyncProcessor:
|
||||||
"""Startup: for each registered handler, fetch config for all its
|
"""Startup: for each registered handler, fetch config for all its
|
||||||
types across all workspaces and invoke the handler once per
|
types across all workspaces and invoke the handler once per
|
||||||
workspace. Retries until successful — config service may not be
|
workspace. Retries until successful — config service may not be
|
||||||
ready yet."""
|
ready yet.
|
||||||
|
|
||||||
|
Tracks which workspaces have already been applied so that
|
||||||
|
retries resume from where they left off rather than re-applying
|
||||||
|
all workspaces from scratch.
|
||||||
|
"""
|
||||||
|
|
||||||
|
applied_ws = set()
|
||||||
|
|
||||||
while self.running:
|
while self.running:
|
||||||
|
|
||||||
|
|
@ -225,11 +243,15 @@ class AsyncProcessor:
|
||||||
for ws, kv in type_data.items():
|
for ws, kv in type_data.items():
|
||||||
per_ws.setdefault(ws, {})[t] = kv
|
per_ws.setdefault(ws, {})[t] = kv
|
||||||
|
|
||||||
# Call the handler once per workspace
|
# Call the handler once per workspace, skipping
|
||||||
|
# workspaces already applied in a previous attempt
|
||||||
for ws, config in per_ws.items():
|
for ws, config in per_ws.items():
|
||||||
if ws.startswith("_"):
|
if ws.startswith("_"):
|
||||||
continue
|
continue
|
||||||
|
if ws in applied_ws:
|
||||||
|
continue
|
||||||
await entry["handler"](ws, config, version)
|
await entry["handler"](ws, config, version)
|
||||||
|
applied_ws.add(ws)
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
f"Applied startup config version {version}"
|
f"Applied startup config version {version}"
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue