Skip to content

Commit d71019a

Browse files
authored
Merge pull request #1250 from RoboFinSystems/bugfix/parquet-path-reads
Sweep the remaining bare-path pyarrow reads; bump the demo client floor
2 parents 1c1bc77 + 461d2ff commit d71019a

7 files changed

Lines changed: 125 additions & 65 deletions

File tree

‎.gitignore‎

Lines changed: 1 addition & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -147,10 +147,5 @@ robosystems/adapters/sec/arelle/cache/
147147

148148
# dbt per-machine state (anonymous usage id)
149149
**/dbt/.user.yml
150-
# Collateral of the line above: ArelleClient._setup_cache_directory sets the
151-
# process-global XDG_CACHE_HOME to robosystems/adapters/sec so Arelle resolves
152-
# its own cache, and every other XDG-respecting library in that process then
153-
# caches there too. huggingface_hub is the one that shows up locally
154-
# (.agent_harnesses.json + xet/logs). Ignore the whole tree, not one file —
155-
# the next such library will land here as well.
150+
# HF cache — ArelleClient points XDG_CACHE_HOME at this directory's parent
156151
robosystems/adapters/sec/huggingface/

‎pyproject.toml‎

Lines changed: 5 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -37,9 +37,6 @@ dependencies = [
3737
"sse-starlette>=3.3.0,<4.0",
3838
"strawberry-graphql[fastapi]>=0.314.0,<1.0",
3939
"uvicorn>=0.48.0,<1.0",
40-
# uvloop ships no Windows wheels and its sdist refuses to build there, so an
41-
# unconditional dependency makes `uv sync` fail outright on Windows hosts.
42-
# Mirrors the marker uvicorn[standard] already carries for its own uvloop dep.
4340
"uvloop>=0.21.0,<1.0; sys_platform != 'win32'",
4441
# Cloud Service (AWS)
4542
"boto3>=1.42.0,<2.0",
@@ -55,9 +52,7 @@ dependencies = [
5552
"redis>=8.0.1,<9.0",
5653
# Document Search (OpenSearch)
5754
"opensearch-py>=3.0.0,<4.0",
58-
# Data Orchestration (Dagster) — the dagster family is version-coupled
59-
# (core/webserver on 1.x.y, libraries on the matching 0.(x+16).y);
60-
# bump all floors together, never one package at a time
55+
# Data Orchestration (Dagster)
6156
"dagster>=1.13.15,<2.0",
6257
"dagster-webserver>=1.13.15,<2.0",
6358
"dagster-postgres>=0.29.15,<1.0",
@@ -77,9 +72,7 @@ dependencies = [
7772
"intuit-oauth>=1.2.0,<2.0",
7873
"python-quickbooks>=0.9.0,<1.0",
7974
"stripe>=11.1.0,<12.0",
80-
# Observability — the opentelemetry-* family versions are coupled
81-
# (api/sdk/exporters track 1.x, instrumentations track the matching 0.xxb0);
82-
# bump all floors together, never one package at a time
75+
# Observability
8376
"opentelemetry-api>=1.44.0,<2.0",
8477
"opentelemetry-exporter-otlp>=1.44.0,<2.0",
8578
"opentelemetry-sdk>=1.44.0,<2.0",
@@ -89,13 +82,9 @@ dependencies = [
8982
"opentelemetry-instrumentation-requests>=0.65b0,<1.0",
9083
# Utilities
9184
"aiohttp>=3.14.3",
92-
"authlib>=1.6.5,<2.0", # OIDC login: PKCE + token-endpoint client (id_token validation rides pyjwt)
85+
"authlib>=1.6.5,<2.0",
9386
"bcrypt>=5.0.0,<6.0",
9487
"beautifulsoup4>=4.13.0,<5.0",
95-
# Direct dependency: auth cache, credential encryption, and connection
96-
# credentials import it at module level. Without this pin it reaches the
97-
# prod closure only transitively (intuit-oauth → pyjwt[crypto]), so
98-
# dropping that adapter would break the auth path at import time.
9988
"cryptography>=50.0.0,<51.0",
10089
"email-validator>=2.2.0,<3.0",
10190
"holidays>=0.101,<1.0",
@@ -113,19 +102,14 @@ dependencies = [
113102
"requests>=2.34.2,<3.0",
114103
"retrying>=1.4.0,<2.0",
115104
"uuid6>=2025.0.0",
116-
# Passkey MFA: RP-side WebAuthn ceremony generation/verification
117-
# (py_webauthn; cbor2 rides transitively)
118105
"webauthn>=2.7.0,<3.0",
119106
"pyshacl>=0.26,<0.31",
120107
]
121108

122109
[project.optional-dependencies]
123110
dev = [
124-
# RoboSystems client for demo and testing. Ranged, not pinned: the demos
125-
# are the dogfooding surface and should exercise the current client, and
126-
# uv.lock is what makes a given checkout reproducible. <2 is the SDK's
127-
# post-1.0 semver boundary.
128-
"robosystems-client>=1.8.1,<2.0",
111+
# RoboSystems client for demo and testing.
112+
"robosystems-client>=1.11.0,<2.0",
129113

130114
# Testing framework
131115
"pytest>=9.0.0,<10.0",
@@ -258,18 +242,6 @@ skip-magic-trailing-comma = false
258242
line-ending = "auto"
259243

260244
[tool.basedpyright]
261-
# Include every Python-bearing dir that ships code:
262-
# - ``robosystems`` and ``main.py``: core platform
263-
# - ``examples``: demo scripts (good-faith use of real SDK paths)
264-
# - ``bin``: Lambda functions (rotation handlers, volume monitors —
265-
# production code, type-check is load-bearing)
266-
# The project's rule suppressions below cover patterns demo and infra
267-
# scripts use (untyped SDK / boto3 responses); inclusion fixes the IDE
268-
# experience without adding CI noise.
269-
# ``tests`` stays excluded — 600+ files of mock-heavy code that produce
270-
# manageable runtime correctness via the test suite itself, and the
271-
# basedpyright noise outweighs the catch rate. ``migrations`` stays
272-
# excluded (autogenerated SQL).
273245
include = ["robosystems", "main.py", "examples", "bin"]
274246
exclude = ["**/__pycache__", ".venv", "robosystems/adapters/sec/arelle", "tests"]
275247
extraPaths = ["."]

‎robosystems/adapters/sec/enrichment.py‎

Lines changed: 15 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -360,7 +360,8 @@ def _load_element_knowledge(self) -> dict[str, dict] | None:
360360

361361
import pyarrow.parquet as pq
362362

363-
table = pq.read_table(path)
363+
with open(path, "rb") as f:
364+
table = pq.read_table(f)
364365
result = {}
365366

366367
# Read columnar and transpose to row dicts
@@ -383,7 +384,7 @@ def _load_element_knowledge(self) -> dict[str, dict] | None:
383384
logger.info(f"Loaded element knowledge artifact: {len(result)} elements")
384385
return result
385386
except Exception as e:
386-
logger.debug(f"Element knowledge artifact not available: {e}")
387+
logger.warning(f"Element knowledge artifact failed to load: {e}")
387388
return None
388389

389390
def _load_structure_profiles(self) -> dict[str, dict[str, float]] | None:
@@ -396,7 +397,8 @@ def _load_structure_profiles(self) -> dict[str, dict[str, float]] | None:
396397

397398
import pyarrow.parquet as pq
398399

399-
table = pq.read_table(path)
400+
with open(path, "rb") as f:
401+
table = pq.read_table(f)
400402
columns = table.to_pydict()
401403

402404
result: dict[str, dict[str, float]] = {}
@@ -408,7 +410,7 @@ def _load_structure_profiles(self) -> dict[str, dict[str, float]] | None:
408410
logger.info(f"Loaded structure profiles artifact: {len(result)} types")
409411
return result
410412
except Exception as e:
411-
logger.debug(f"Structure profiles artifact not available: {e}")
413+
logger.warning(f"Structure profiles artifact failed to load: {e}")
412414
return None
413415

414416
def _load_structure_consensus(self) -> dict[str, dict] | None:
@@ -421,7 +423,8 @@ def _load_structure_consensus(self) -> dict[str, dict] | None:
421423

422424
import pyarrow.parquet as pq
423425

424-
table = pq.read_table(path)
426+
with open(path, "rb") as f:
427+
table = pq.read_table(f)
425428
columns = table.to_pydict()
426429

427430
result = {}
@@ -435,7 +438,7 @@ def _load_structure_consensus(self) -> dict[str, dict] | None:
435438
logger.info(f"Loaded structure consensus artifact: {len(result)} entries")
436439
return result
437440
except Exception as e:
438-
logger.debug(f"Structure consensus artifact not available: {e}")
441+
logger.warning(f"Structure consensus artifact failed to load: {e}")
439442
return None
440443

441444
@property
@@ -462,7 +465,8 @@ def _load_disclosure_profiles(self) -> dict[str, dict[str, float]] | None:
462465

463466
import pyarrow.parquet as pq
464467

465-
table = pq.read_table(path)
468+
with open(path, "rb") as f:
469+
table = pq.read_table(f)
466470
columns = table.to_pydict()
467471

468472
result: dict[str, dict[str, float]] = {}
@@ -474,7 +478,7 @@ def _load_disclosure_profiles(self) -> dict[str, dict[str, float]] | None:
474478
logger.info(f"Loaded disclosure profiles artifact: {len(result)} types")
475479
return result
476480
except Exception as e:
477-
logger.debug(f"Disclosure profiles artifact not available: {e}")
481+
logger.warning(f"Disclosure profiles artifact failed to load: {e}")
478482
return None
479483

480484
def _load_disclosure_consensus(self) -> dict[str, dict] | None:
@@ -487,7 +491,8 @@ def _load_disclosure_consensus(self) -> dict[str, dict] | None:
487491

488492
import pyarrow.parquet as pq
489493

490-
table = pq.read_table(path)
494+
with open(path, "rb") as f:
495+
table = pq.read_table(f)
491496
columns = table.to_pydict()
492497

493498
result = {}
@@ -501,7 +506,7 @@ def _load_disclosure_consensus(self) -> dict[str, dict] | None:
501506
logger.info(f"Loaded disclosure consensus artifact: {len(result)} entries")
502507
return result
503508
except Exception as e:
504-
logger.debug(f"Disclosure consensus artifact not available: {e}")
509+
logger.warning(f"Disclosure consensus artifact failed to load: {e}")
505510
return None
506511

507512
# -- Embedding ------------------------------------------------------------

‎robosystems/adapters/sec/knowledge/artifact.py‎

Lines changed: 2 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -581,23 +581,15 @@ def _load_pagerank_scores(self) -> dict[str, float]:
581581
return {}
582582

583583
try:
584-
# Read through a handle. This module imports networkit at import time, and
585-
# in a process that has done so `pq.read_table(path)` raises ArrowKeyError
586-
# from the LocalFileSystem constructor (icebug bundles a second libarrow
587-
# that collides with pyarrow's on Arrow's global filesystem registry — see
588-
# this package's __init__). Passing a path here meant this function could
589-
# never succeed: it raised on every call, the except below swallowed it at
590-
# debug level, and disclosure profiles silently fell back to frequency-only
591-
# weighting. The writers in this file already take the handle path.
584+
# Handle, not a path — see this package's __init__.
592585
with open(path, "rb") as f:
593586
table = pq.read_table(f, columns=["qname", "pagerank"])
594587
cols = table.to_pydict()
595588
return {
596589
cols["qname"][i]: cols["pagerank"][i] or 0.0 for i in range(len(cols["qname"]))
597590
}
598591
except Exception as e:
599-
# Warning, not debug: the empty dict degrades weighting silently, so a
600-
# recurring failure here has to be visible in logs to be noticed at all.
592+
# Warning, not debug: the empty dict degrades weighting silently.
601593
logger.warning(f"Failed to load PageRank scores from {path}: {e}")
602594
return {}
603595

‎robosystems/adapters/sec/processors/consolidation.py‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -154,10 +154,12 @@ def consolidate_parquet_from_disk(
154154
tables = []
155155
for pq_file in parquet_files:
156156
try:
157-
table = pq.read_table(pq_file)
157+
# Handle, not a path — see adapters/sec/knowledge/__init__.py.
158+
with open(pq_file, "rb") as f:
159+
table = pq.read_table(f)
158160
tables.append(table)
159161
except Exception as e:
160-
logger.warning("Skipping corrupted parquet file %s: %s", pq_file, e)
162+
logger.warning("Skipping unreadable parquet file %s: %s", pq_file, e)
161163
continue
162164

163165
if not tables:
Lines changed: 94 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,94 @@
1+
"""Artifact loaders must work in a process that has imported networkit.
2+
3+
icebug bundles a second libarrow beside pyarrow's; the two collide on Arrow's
4+
global filesystem registry, so once networkit is loaded `pq.read_table(path)`
5+
raises ArrowKeyError forever. The knowledge builders import networkit, so any
6+
loader reading their output has to take a file handle. This test poisons the
7+
registry the same way and asserts the loaders still return data — the previous
8+
tests all mocked the read, so a bare path passed them and failed in production.
9+
"""
10+
11+
from __future__ import annotations
12+
13+
from pathlib import Path
14+
15+
import networkit # noqa: F401 — poisons the Arrow registry, exactly as prod does
16+
import pyarrow as pa
17+
import pyarrow.parquet as pq
18+
import pytest
19+
20+
from robosystems.adapters.sec.enrichment import SemanticEnricher
21+
22+
pytestmark = pytest.mark.unit
23+
24+
25+
def _write(path: Path, columns: dict) -> None:
26+
with open(path, "wb") as f:
27+
pq.write_table(pa.table(columns), f)
28+
29+
30+
@pytest.fixture
31+
def artifacts(tmp_path, monkeypatch):
32+
_write(
33+
tmp_path / "element_knowledge.parquet",
34+
{
35+
"qname": ["us-gaap:Assets"],
36+
"primary_statement": ["BalanceSheet"],
37+
"bfs_depth": [1],
38+
"pagerank": [0.9],
39+
"core_number": [2],
40+
"neighborhood_agreement": [0.5],
41+
"filing_count": [3],
42+
"disclosure_type": ["x"],
43+
},
44+
)
45+
_write(
46+
tmp_path / "structure_profiles.parquet",
47+
{
48+
"canonical_type": ["BalanceSheet"],
49+
"qname": ["us-gaap:Assets"],
50+
"frequency": [0.8],
51+
"structure_count": [4],
52+
},
53+
)
54+
monkeypatch.setattr(
55+
"robosystems.config.storage.shared.get_artifact_path",
56+
lambda name: str(tmp_path / f"{name}.parquet"),
57+
)
58+
return tmp_path
59+
60+
61+
def _enricher() -> SemanticEnricher:
62+
# __new__, not __init__: construction loads an embedding model, and the read
63+
# path under test does not need one.
64+
return SemanticEnricher.__new__(SemanticEnricher)
65+
66+
67+
def test_element_knowledge_loads(artifacts):
68+
result = _enricher()._load_element_knowledge()
69+
assert result is not None, "bare-path read — ArrowKeyError swallowed to None"
70+
assert result["us-gaap:Assets"]["pagerank"] == 0.9
71+
72+
73+
def test_structure_profiles_load(artifacts):
74+
result = _enricher()._load_structure_profiles()
75+
assert result is not None, "bare-path read — ArrowKeyError swallowed to None"
76+
assert result["BalanceSheet"]["us-gaap:Assets"] == pytest.approx(0.8)
77+
78+
79+
def test_consolidation_reads_parquet_from_disk(tmp_path):
80+
"""The consolidation read attributes any failure to file corruption, so a
81+
registry collision there discards real rows under a misleading label."""
82+
from robosystems.adapters.sec.processors.consolidation import (
83+
consolidate_parquet_from_disk,
84+
)
85+
86+
table_dir = tmp_path / "nodes/Entity"
87+
table_dir.mkdir(parents=True)
88+
_write(table_dir / "a.parquet", {"identifier": ["e1"], "name": ["Acme"]})
89+
90+
out = consolidate_parquet_from_disk(tmp_path, "nodes/Entity")
91+
assert out is not None, "every file read as corrupt — rows silently dropped"
92+
from io import BytesIO
93+
94+
assert pq.read_table(BytesIO(out)).num_rows == 1

‎uv.lock‎

Lines changed: 4 additions & 4 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

0 commit comments

Comments
 (0)