Skip to content

Commit e678986

Browse files
satraclaude
andcommitted
feat(040): Parquet-only commit + embeddings at commit + cross-ref resolution
T011: commit reads from ParquetStore, writes to ParquetStore (no YAML output) T012: _resolve_cross_references uses ParquetStore.list() instead of YAML globs T013: Removed yaml.dump from commit path T017: Embedding computation added to commit — batch compute after sha256 T018: All entity types (elements, schemas, values, valuesets) get embeddings Fix: ParquetStore dedup skips empty sha256 (pre-commit staged entities). 28 tests pass (5 commit + 23 parquet). 21/33 tasks done. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent 41faf32 commit e678986

4 files changed

Lines changed: 235 additions & 226 deletions

File tree

library/src/undata_library/commit.py

Lines changed: 127 additions & 138 deletions
Original file line numberDiff line numberDiff line change
@@ -9,8 +9,8 @@
99

1010
import yaml
1111

12-
from .hashing import compute_identity_hash, determine_hash_mode, generate_short_key
13-
from .utils import safe_load_yaml, sanitize_filename
12+
from .hashing import compute_identity_hash, determine_hash_mode
13+
from .utils import safe_load_yaml
1414

1515
logger = logging.getLogger(__name__)
1616

@@ -45,10 +45,12 @@ def commit_staged(
4545
if output_dir is None and output_backend is not None and hasattr(output_backend, "base_dir"):
4646
output_dir = output_backend.base_dir
4747
stats = {"committed": 0, "merged": 0, "rejected": 0, "per_type": {}}
48-
# Track staging filename → sha256 for flag entity_ref resolution
4948
staging_to_sha256: dict[str, str] = {}
5049

5150
from .staging import iter_staged
51+
from .storage.parquet_store import ParquetStore
52+
53+
pq_registry = ParquetStore(output_dir)
5254

5355
for entity_type in ("elements", "schemas", "values", "valuesets"):
5456
type_dir = staging_dir / entity_type
@@ -106,59 +108,60 @@ def commit_staged(
106108
ontology_anchored=ontology_anchored,
107109
primary_ontology_uri=primary_uri,
108110
)
109-
key = generate_short_key(sha256)
110-
111111
# Record staging identifier → sha256 mapping for flag resolution
112112
staging_id = data.get("_identifier", data.get("file_name", ""))
113113
if staging_id:
114114
staging_to_sha256[staging_id] = sha256
115115
if "." in staging_id:
116116
staging_to_sha256[staging_id.rsplit(".", 1)[0]] = sha256
117117

118-
# Check if any existing file has this hash (for cross-source merge)
119-
existing_with_hash = list(out_dir.glob(f"*_{key}.yaml"))
120-
if existing_with_hash:
121-
target = existing_with_hash[0]
122-
else:
123-
name = _derive_name(data, entity_type)
124-
safe_name = sanitize_filename(name)
125-
target = out_dir / f"{safe_name}_{key}.yaml"
126-
127-
if target.exists():
128-
# Merge provenance
129-
_merge_provenance(target, data)
118+
# Check if entity already exists in registry (for cross-source merge)
119+
existing = pq_registry.read(entity_type, sha256)
120+
if existing:
121+
# Merge provenance into existing
122+
existing_prov = existing.get("provenance", [])
123+
existing_keys = {
124+
(p.get("source", ""), p.get("name", ""))
125+
for p in existing_prov
126+
if isinstance(p, dict)
127+
}
128+
for p in provenance:
129+
pk = (p.get("source", ""), p.get("name", ""))
130+
if isinstance(p, dict) and pk not in existing_keys:
131+
existing_prov.append(p)
132+
data["provenance"] = existing_prov
130133
type_merged += 1
131134
else:
132-
# Write new file with sha256
133-
data["sha256"] = sha256
134-
target.write_text(
135-
yaml.dump(data, default_flow_style=False, sort_keys=False),
136-
encoding="utf-8",
137-
)
138135
type_committed += 1
139136

140-
# Also accumulate for Parquet batch write
141137
data["sha256"] = sha256
138+
data["file_name"] = _derive_name(data, entity_type)
142139
committed_entities.append(data)
143140

144141
stats["per_type"][entity_type] = {"committed": type_committed, "merged": type_merged}
145142
stats["committed"] += type_committed
146143
stats["merged"] += type_merged
147144

148-
# Write committed entities as Parquet (in addition to YAML for now)
145+
# Compute embeddings for all committed entities in batch
149146
if committed_entities:
150-
from .storage.parquet_store import ParquetStore
147+
try:
148+
from .embeddings import compute_entity_embeddings
151149

152-
pq_store = ParquetStore(output_dir)
150+
committed_entities = compute_entity_embeddings(committed_entities)
151+
except ImportError:
152+
logger.debug("sentence-transformers not available; skipping embeddings")
153+
154+
# Write committed entities to Parquet (sole output format)
155+
if committed_entities:
153156
source = ""
154157
if committed_entities[0].get("provenance"):
155158
prov = committed_entities[0]["provenance"]
156159
if prov and isinstance(prov[0], dict):
157160
source = prov[0].get("source", "committed")
158-
pq_store.write_batch(entity_type, committed_entities, source=source or "committed")
161+
pq_registry.write_batch(entity_type, committed_entities, source=source or "committed")
159162

160-
# Post-commit: resolve schema properties and valueset members to sha256 hashes
161-
_resolve_cross_references(output_dir)
163+
# Post-commit: resolve cross-references using ParquetStore
164+
_resolve_cross_references(output_dir, pq_registry)
162165

163166
# Post-commit: resolve curation flag entity_refs from filenames to sha256 hashes
164167
_resolve_flag_entity_refs(output_dir, staging_to_sha256)
@@ -170,77 +173,61 @@ def commit_staged(
170173
return stats
171174

172175

173-
def _resolve_cross_references(output_dir: Path) -> None:
174-
"""Resolve schema properties and valueset members from names to sha256 hashes.
175-
176-
After all entities are committed with their sha256 hashes, this step updates:
177-
- Schema properties: slot names → element sha256 hashes
178-
- Valueset members: value labels → value sha256 hashes
176+
def _resolve_cross_references(output_dir: Path, pq_store=None) -> None:
177+
"""Resolve schema properties, valueset members, and type_refs via ParquetStore.
179178
180-
Uses a name→sha256 lookup built from committed elements and values.
179+
Reads all entities from ParquetStore, builds lookup dicts, resolves references,
180+
and writes back resolved entities.
181181
"""
182-
# Build element (class, name) → sha256 lookup for class-aware resolution
183-
# Key: (class_name, slot_name) → sha256
182+
from .storage.parquet_store import ParquetStore
183+
184+
store = pq_store or ParquetStore(output_dir)
185+
186+
# Build element (class, name) → sha256 lookup
184187
elem_by_class: dict[tuple[str, str], str] = {}
185-
# Fallback: name → sha256 (first encountered)
186188
elem_by_name: dict[str, str] = {}
187-
elements_dir = output_dir / "elements"
188-
if elements_dir.exists():
189-
for f in elements_dir.glob("*.yaml"):
190-
data = safe_load_yaml(f)
191-
if not data or "sha256" not in data:
192-
continue
193-
sha = data["sha256"]
194-
for prov in data.get("provenance", []):
195-
if isinstance(prov, dict) and prov.get("name"):
196-
name = prov["name"]
197-
cls = prov.get("class", prov.get("class_", ""))
198-
if cls:
199-
elem_by_class[(cls, name)] = sha
200-
elem_by_class[(cls, name.lower())] = sha
201-
if name not in elem_by_name:
202-
elem_by_name[name] = sha
203-
if name.lower() not in elem_by_name:
204-
elem_by_name[name.lower()] = sha
189+
for data in store.list("elements"):
190+
sha = data.get("sha256", "")
191+
if not sha:
192+
continue
193+
for prov in data.get("provenance", []):
194+
if isinstance(prov, dict) and prov.get("name"):
195+
name = prov["name"]
196+
cls = prov.get("class", prov.get("class_", ""))
197+
if cls:
198+
elem_by_class[(cls, name)] = sha
199+
elem_by_class[(cls, name.lower())] = sha
200+
if name not in elem_by_name:
201+
elem_by_name[name] = sha
202+
if name.lower() not in elem_by_name:
203+
elem_by_name[name.lower()] = sha
205204

206205
# Build value label → sha256 lookup
207206
val_lookup: dict[str, str] = {}
208-
values_dir = output_dir / "values"
209-
if values_dir.exists():
210-
for f in values_dir.glob("*.yaml"):
211-
data = safe_load_yaml(f)
212-
if not data or "sha256" not in data:
213-
continue
214-
sha = data["sha256"]
215-
sem = data.get("semantic", {})
216-
label = sem.get("label", "")
217-
if label:
218-
if label not in val_lookup:
219-
val_lookup[label] = sha
220-
if label.lower() not in val_lookup:
221-
val_lookup[label.lower()] = sha
222-
# Also index by provenance name
223-
for prov in data.get("provenance", []):
224-
if isinstance(prov, dict) and prov.get("name"):
225-
name = prov["name"]
226-
if name not in val_lookup:
227-
val_lookup[name] = sha
228-
if name.lower() not in val_lookup:
229-
val_lookup[name.lower()] = sha
230-
231-
# Update schema properties — class-aware resolution
232-
schemas_dir = output_dir / "schemas"
233-
if schemas_dir.exists() and (elem_by_class or elem_by_name):
234-
for f in schemas_dir.glob("*.yaml"):
235-
data = safe_load_yaml(f)
236-
if not data:
237-
continue
207+
for data in store.list("values"):
208+
sha = data.get("sha256", "")
209+
if not sha:
210+
continue
211+
label = data.get("label", data.get("semantic", {}).get("label", ""))
212+
if label:
213+
val_lookup.setdefault(label, sha)
214+
val_lookup.setdefault(label.lower(), sha)
215+
for prov in data.get("provenance", []):
216+
if isinstance(prov, dict) and prov.get("name"):
217+
name = prov["name"]
218+
val_lookup.setdefault(name, sha)
219+
val_lookup.setdefault(name.lower(), sha)
220+
221+
# Update schema properties — class-aware resolution via ParquetStore
222+
if elem_by_class or elem_by_name:
223+
resolved_schemas = []
224+
for data in store.list("schemas"):
238225
sem = data.get("semantic", {})
239226
props = sem.get("properties", [])
240227
if not props:
228+
resolved_schemas.append(data)
241229
continue
242230

243-
# Get the schema's class name from provenance
244231
schema_class = ""
245232
for prov in data.get("provenance", []):
246233
if isinstance(prov, dict):
@@ -250,11 +237,9 @@ def _resolve_cross_references(output_dir: Path) -> None:
250237
resolved = []
251238
changed = False
252239
for prop in props:
253-
# If it's already a sha256 hash (64 hex chars), keep it
254240
if len(prop) == 64 and all(c in "0123456789abcdef" for c in prop):
255241
resolved.append(prop)
256242
continue
257-
# Class-aware: prefer element from same class as the schema
258243
sha = (
259244
elem_by_class.get((schema_class, prop))
260245
or elem_by_class.get((schema_class, prop.lower()))
@@ -265,24 +250,24 @@ def _resolve_cross_references(output_dir: Path) -> None:
265250
resolved.append(sha)
266251
changed = True
267252
else:
268-
resolved.append(prop) # Keep unresolved name
253+
resolved.append(prop)
269254
if changed:
270255
sem["properties"] = resolved
271-
f.write_text(
272-
yaml.dump(data, default_flow_style=False, sort_keys=False),
273-
encoding="utf-8",
274-
)
256+
data["semantic"] = sem
257+
resolved_schemas.append(data)
275258

276-
# Update valueset members
277-
valuesets_dir = output_dir / "valuesets"
278-
if valuesets_dir.exists() and val_lookup:
279-
for f in valuesets_dir.glob("*.yaml"):
280-
data = safe_load_yaml(f)
281-
if not data:
282-
continue
259+
if resolved_schemas:
260+
source = _get_source(resolved_schemas)
261+
store.write_batch("schemas", resolved_schemas, source=source)
262+
263+
# Update valueset members via ParquetStore
264+
if val_lookup:
265+
resolved_vs = []
266+
for data in store.list("valuesets"):
283267
sem = data.get("semantic", {})
284268
members = sem.get("members", [])
285269
if not members:
270+
resolved_vs.append(data)
286271
continue
287272
resolved = []
288273
changed = False
@@ -298,46 +283,50 @@ def _resolve_cross_references(output_dir: Path) -> None:
298283
resolved.append(member)
299284
if changed:
300285
sem["members"] = resolved
301-
f.write_text(
302-
yaml.dump(data, default_flow_style=False, sort_keys=False),
303-
encoding="utf-8",
304-
)
286+
data["semantic"] = sem
287+
resolved_vs.append(data)
288+
289+
if resolved_vs:
290+
source = _get_source(resolved_vs)
291+
store.write_batch("valuesets", resolved_vs, source=source)
305292

306293
# Resolve element type_ref: class names → schema sha256 hashes
307-
# Build schema name → sha256 lookup
308294
schema_lookup: dict[str, str] = {}
309-
if schemas_dir.exists():
310-
for f in schemas_dir.glob("*.yaml"):
311-
data = safe_load_yaml(f)
312-
if not data or "sha256" not in data:
313-
continue
314-
sha = data["sha256"]
315-
for prov in data.get("provenance", []):
316-
if isinstance(prov, dict):
317-
name = prov.get("name", prov.get("class", ""))
318-
if name:
319-
schema_lookup[name] = sha
320-
schema_lookup[name.lower()] = sha
321-
322-
if elements_dir.exists() and schema_lookup:
323-
for f in elements_dir.glob("*.yaml"):
324-
data = safe_load_yaml(f)
325-
if not data:
326-
continue
295+
for data in store.list("schemas"):
296+
sha = data.get("sha256", "")
297+
if not sha:
298+
continue
299+
for prov in data.get("provenance", []):
300+
if isinstance(prov, dict):
301+
name = prov.get("name", prov.get("class", ""))
302+
if name:
303+
schema_lookup[name] = sha
304+
schema_lookup[name.lower()] = sha
305+
306+
if schema_lookup:
307+
resolved_elems = []
308+
for data in store.list("elements"):
327309
sem = data.get("semantic", {})
328310
type_ref = sem.get("type_ref")
329-
if not type_ref:
330-
continue
331-
# Already a sha256?
332-
if len(type_ref) == 64 and all(c in "0123456789abcdef" for c in type_ref):
333-
continue
334-
sha = schema_lookup.get(type_ref) or schema_lookup.get(type_ref.lower())
335-
if sha:
336-
sem["type_ref"] = sha
337-
f.write_text(
338-
yaml.dump(data, default_flow_style=False, sort_keys=False),
339-
encoding="utf-8",
340-
)
311+
if type_ref and len(type_ref) != 64:
312+
sha = schema_lookup.get(type_ref) or schema_lookup.get(type_ref.lower())
313+
if sha:
314+
sem["type_ref"] = sha
315+
data["semantic"] = sem
316+
resolved_elems.append(data)
317+
318+
if resolved_elems:
319+
source = _get_source(resolved_elems)
320+
store.write_batch("elements", resolved_elems, source=source)
321+
322+
323+
def _get_source(entities: list[dict]) -> str:
324+
"""Extract source name from first entity's provenance."""
325+
if entities and entities[0].get("provenance"):
326+
prov = entities[0]["provenance"]
327+
if prov and isinstance(prov[0], dict):
328+
return prov[0].get("source", "committed")
329+
return "committed"
341330

342331

343332
def _resolve_flag_entity_refs(

library/src/undata_library/storage/parquet_store.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -175,7 +175,7 @@ def write_batch(
175175
serialized = _serialize_entity(entity, source)
176176
sha = serialized["sha256"]
177177

178-
if sha in existing:
178+
if sha and sha in existing:
179179
# Merge provenance
180180
old_prov = json.loads(existing[sha].get("provenance", "[]"))
181181
new_prov = json.loads(serialized["provenance"])
@@ -186,7 +186,9 @@ def write_batch(
186186
serialized["provenance"] = json.dumps(old_prov, default=str)
187187

188188
rows.append(serialized)
189-
existing[sha] = serialized # update for subsequent dedup within batch
189+
# Use sha as dedup key, or unique index for pre-commit entities without sha
190+
dedup_key = sha if sha else f"__pending_{len(rows)}__"
191+
existing[dedup_key] = serialized
190192

191193
# Write all (existing merged + new)
192194
all_rows = list(existing.values())

0 commit comments

Comments
 (0)