Skip to content

Commit bb5f4b0

Browse files
committed
feat(store): stage validate=true registration behind a committed overlay
Add a staging overlay to GtsStore: staged entities are visible to internal validation and the $ref reference registry but invisible to public reads (get_committed) until commit. A validate=true registration stages, validates, then commits on success or discards on failure, so a concurrent reader never observes an unvalidated entity and no lock is held across validation. Batch add_schemas now runs two-phase - stage every entry, validate each against the fully-staged set, then commit the survivors and discard the rest - making it order-independent while never publishing an entry that fails. Add a threaded concurrency test. Signed-off-by: Artifizer <artifizer@gmail.com>
1 parent df21598 commit bb5f4b0

3 files changed

Lines changed: 317 additions & 78 deletions

File tree

‎gts/src/gts/ops.py‎

Lines changed: 166 additions & 74 deletions
Original file line numberDiff line numberDiff line change
@@ -444,37 +444,40 @@ def add_entity(
444444
else:
445445
assert entity.raw_id is not None
446446
store_key = entity.raw_id
447-
with self.store.transaction():
448-
previous = self.store.get(store_key)
449-
if (
450-
previous
451-
and not self.allow_entity_updates
452-
and previous.content != entity.content
453-
):
454-
return GtsAddEntityResult(
455-
ok=False,
456-
error=f"Entity '{store_key}' is already registered with different content",
457-
is_type_schema=entity.is_schema,
458-
conflict=True,
459-
)
460-
self.store.register(entity)
447+
previous = self.store.get_committed(store_key)
448+
if (
449+
previous
450+
and not self.allow_entity_updates
451+
and previous.content != entity.content
452+
):
453+
return GtsAddEntityResult(
454+
ok=False,
455+
error=f"Entity '{store_key}' is already registered with different content",
456+
is_type_schema=entity.is_schema,
457+
conflict=True,
458+
)
461459

462-
try:
463-
if entity.is_schema:
464-
self.store.validate_schema_basic(store_key)
465-
if validate:
466-
self.store.validate_schema(store_key, gts_ref_validation)
467-
elif validate:
468-
self.store.validate_instance(store_key, gts_ref_validation)
469-
except Exception as e: # noqa: BLE001 - converted to a result object at API boundary
470-
self.store.unregister(store_key)
471-
if previous:
472-
self.store.register(previous)
473-
return GtsAddEntityResult(
474-
ok=False,
475-
error=f"Validation failed: {e!s}",
476-
is_type_schema=entity.is_schema,
477-
)
460+
# Stage the entity (invisible to public reads) and validate it before
461+
# publishing. A failure discards the staged copy, so a reader never
462+
# observes an entity that has not passed validation, and the committed
463+
# state (any prior version under this id) is never touched. No lock is
464+
# held across validation, so concurrent reads are not blocked.
465+
self.store.stage(entity)
466+
try:
467+
if entity.is_schema:
468+
self.store.validate_schema_basic(store_key)
469+
if validate:
470+
self.store.validate_schema(store_key, gts_ref_validation)
471+
elif validate:
472+
self.store.validate_instance(store_key, gts_ref_validation)
473+
except Exception as e: # noqa: BLE001 - converted to a result object at API boundary
474+
self.store.discard(store_key)
475+
return GtsAddEntityResult(
476+
ok=False,
477+
error=f"Validation failed: {e!s}",
478+
is_type_schema=entity.is_schema,
479+
)
480+
self.store.commit(store_key)
478481

479482
# Return gts_id if available, otherwise raw_id
480483
entity_id = entity.gts_id.id if entity.gts_id else (entity.raw_id or "")
@@ -494,39 +497,33 @@ def add_entities(
494497
ok = all(r.ok for r in results)
495498
return GtsAddEntitiesResult(ok=ok, results=results)
496499

497-
def add_schemas(
500+
def _validate_staged(
498501
self,
499-
schemas: builtins.list[dict[str, Any]],
500-
validate: bool = False,
501-
gts_ref_validation: GtsRefValidationMode = GtsRefValidationMode.ANY_VALID,
502-
) -> GtsAddSchemasResult:
503-
"""Register a batch of GTS Type Schemas.
504-
505-
Each entry's GTS Type Identifier is derived from its embedded ``$id``;
506-
the aggregate ``ok`` is ``True`` only when every entry registered.
507-
``validate`` / ``gts_ref_validation`` apply to every entry exactly as
508-
they do on ``POST /entities``.
509-
"""
510-
results = [
511-
self.add_schema(schema, validate=validate, gts_ref_validation=gts_ref_validation)
512-
for schema in schemas
513-
]
514-
ok = all(r.ok for r in results)
515-
return GtsAddSchemasResult(ok=ok, results=results)
516-
517-
def add_schema(
518-
self,
519-
schema: dict[str, Any],
520-
validate: bool = False,
521-
gts_ref_validation: GtsRefValidationMode = GtsRefValidationMode.ANY_VALID,
522-
) -> GtsAddSchemaResult:
523-
"""Register a single GTS Type Schema, deriving its type_id from ``$id``.
524-
525-
The embedded ``$schema`` / ``$id`` presence checks are batch-specific;
526-
the actual registration and (when requested) validation reuse the
527-
single-entity :meth:`add_entity` path, so each entry honors ``validate``
528-
/ ``gts_ref_validation`` exactly like a ``POST /entities`` call.
529-
"""
502+
entity: GtsEntity,
503+
store_key: str,
504+
validate: bool,
505+
gts_ref_validation: GtsRefValidationMode,
506+
) -> str | None:
507+
"""Validate a staged entity against the current (staged + committed)
508+
set, returning an error message on failure or ``None`` on success. The
509+
caller stages before and commits/discards after."""
510+
try:
511+
if entity.is_schema:
512+
self.store.validate_schema_basic(store_key)
513+
if validate:
514+
self.store.validate_schema(store_key, gts_ref_validation)
515+
elif validate:
516+
self.store.validate_instance(store_key, gts_ref_validation)
517+
except Exception as e: # noqa: BLE001 - converted to a result object at API boundary
518+
return f"Validation failed: {e!s}"
519+
return None
520+
521+
def _prepare_type_schema(
522+
self, schema: dict[str, Any]
523+
) -> tuple[str, GtsEntity, str] | GtsAddSchemaResult:
524+
"""Run the batch-specific $schema/$id checks and build the entity. On
525+
success returns ``(type_id, entity, store_key)``; on failure returns the
526+
per-item error result. Does not stage or register anything."""
530527
if not isinstance(schema, dict):
531528
return GtsAddSchemaResult(
532529
ok=False,
@@ -558,18 +555,112 @@ def add_schema(
558555
type_id=type_id,
559556
error=f"Invalid GTS Type Schema $id: {embedded_id}",
560557
)
558+
entity = GtsEntity(content=schema, cfg=self.cfg)
559+
if not entity.is_schema or not entity.gts_id:
560+
return GtsAddSchemaResult(
561+
ok=False, type_id=type_id, error="Unable to detect GTS ID in schema"
562+
)
563+
store_key = entity.gts_id.id
564+
previous = self.store.get_committed(store_key)
565+
if (
566+
previous
567+
and not self.allow_entity_updates
568+
and previous.content != entity.content
569+
):
570+
return GtsAddSchemaResult(
571+
ok=False,
572+
type_id=type_id,
573+
error=f"Entity '{store_key}' is already registered with different content",
574+
conflict=True,
575+
)
576+
return (type_id, entity, store_key)
561577

562-
result = self.add_entity(
563-
schema, validate=validate, gts_ref_validation=gts_ref_validation
564-
)
565-
if result.ok:
566-
return GtsAddSchemaResult(ok=True, type_id=type_id)
567-
return GtsAddSchemaResult(
568-
ok=False,
569-
type_id=type_id,
570-
error=result.error,
571-
conflict=result.conflict,
572-
)
578+
def add_schemas(
579+
self,
580+
schemas: builtins.list[dict[str, Any]],
581+
validate: bool = False,
582+
gts_ref_validation: GtsRefValidationMode = GtsRefValidationMode.ANY_VALID,
583+
) -> GtsAddSchemasResult:
584+
"""Register a batch of GTS Type Schemas.
585+
586+
Each entry's GTS Type Identifier is derived from its embedded ``$id``;
587+
the aggregate ``ok`` is ``True`` only when every entry registered.
588+
``validate`` / ``gts_ref_validation`` apply to every entry exactly as
589+
they do on ``POST /entities``.
590+
591+
With ``validate`` the batch runs in two phases so the outcome is
592+
order-independent and nothing invalid is ever published: every
593+
structurally-valid entry is staged first (invisible to public reads),
594+
then each is validated against the fully-staged set - so an entry can
595+
resolve intra-batch references/ancestors regardless of position - and
596+
finally the entries that passed are committed while the rest are
597+
discarded.
598+
"""
599+
gts_ref_validation = _normalize_gts_ref_validation(gts_ref_validation)
600+
if not validate:
601+
results = [
602+
self.add_schema(schema, validate=False, gts_ref_validation=gts_ref_validation)
603+
for schema in schemas
604+
]
605+
return GtsAddSchemasResult(ok=all(r.ok for r in results), results=results)
606+
607+
results: list[GtsAddSchemaResult | None] = [None] * len(schemas)
608+
# Phase 1: stage every structurally-valid entry.
609+
staged: list[tuple[int, str, GtsEntity, str]] = []
610+
for index, schema in enumerate(schemas):
611+
prepared = self._prepare_type_schema(schema)
612+
if isinstance(prepared, GtsAddSchemaResult):
613+
results[index] = prepared
614+
continue
615+
type_id, entity, store_key = prepared
616+
self.store.stage(entity)
617+
staged.append((index, type_id, entity, store_key))
618+
619+
# Phase 2: validate every staged entry against the fully-staged set.
620+
verdicts: list[tuple[int, str, str, str | None]] = []
621+
for index, type_id, entity, store_key in staged:
622+
error = self._validate_staged(entity, store_key, True, gts_ref_validation)
623+
verdicts.append((index, type_id, store_key, error))
624+
625+
# Phase 3: publish the entries that passed, discard the ones that failed.
626+
for index, type_id, store_key, error in verdicts:
627+
if error is None:
628+
self.store.commit(store_key)
629+
results[index] = GtsAddSchemaResult(ok=True, type_id=type_id)
630+
else:
631+
self.store.discard(store_key)
632+
results[index] = GtsAddSchemaResult(
633+
ok=False, type_id=type_id, error=error
634+
)
635+
636+
final = [r for r in results if r is not None]
637+
return GtsAddSchemasResult(ok=all(r.ok for r in final), results=final)
638+
639+
def add_schema(
640+
self,
641+
schema: dict[str, Any],
642+
validate: bool = False,
643+
gts_ref_validation: GtsRefValidationMode = GtsRefValidationMode.ANY_VALID,
644+
) -> GtsAddSchemaResult:
645+
"""Register a single GTS Type Schema, deriving its type_id from ``$id``.
646+
647+
Stages the entity (invisible to public reads), validates it, then
648+
commits on success or discards on failure, so each entry honors
649+
``validate`` / ``gts_ref_validation`` exactly like a ``POST /entities``
650+
call and an invalid schema is never observable.
651+
"""
652+
gts_ref_validation = _normalize_gts_ref_validation(gts_ref_validation)
653+
prepared = self._prepare_type_schema(schema)
654+
if isinstance(prepared, GtsAddSchemaResult):
655+
return prepared
656+
type_id, entity, store_key = prepared
657+
self.store.stage(entity)
658+
error = self._validate_staged(entity, store_key, validate, gts_ref_validation)
659+
if error is not None:
660+
self.store.discard(store_key)
661+
return GtsAddSchemaResult(ok=False, type_id=type_id, error=error)
662+
self.store.commit(store_key)
663+
return GtsAddSchemaResult(ok=True, type_id=type_id)
573664

574665
def validate_id(self, gts_id: str) -> GtsIdValidationResult:
575666
# Check if it's a wildcard pattern (contains *)
@@ -855,7 +946,8 @@ def get_entity(self, gts_id: str) -> GtsGetEntityResult:
855946
GtsGetEntityResult with entity details or error
856947
"""
857948
try:
858-
entity = self.store.get(gts_id)
949+
# Public read: never expose a staged (not-yet-committed) entity.
950+
entity = self.store.get_committed(gts_id)
859951
if not entity:
860952
return GtsGetEntityResult(
861953
ok=False, error=f"Entity '{gts_id}' not found"

‎gts/src/gts/store.py‎

Lines changed: 73 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -139,6 +139,14 @@ def __init__(self, reader: GtsReader | None = None) -> None:
139139
reader: GtsReader instance to populate entities from
140140
"""
141141
self._by_id: dict[str, GtsEntity] = {}
142+
# Entities registered but not yet published. A staged entity is visible
143+
# to internal validation/resolution (via `get`, the reference registry
144+
# and x-gts-ref existence checks) so a batch can resolve intra-batch
145+
# references regardless of entry order, but it is invisible to public
146+
# reads (`get_committed`, `values`, `keys`, `items`, `query`) until
147+
# `commit`. This is how a validate=true registration avoids ever
148+
# exposing an entity that has not yet passed validation.
149+
self._staged: dict[str, GtsEntity] = {}
142150
self._reader = reader
143151
self._lock = threading.RLock()
144152
self._reference_registry: Registry | None = None
@@ -228,15 +236,31 @@ def register_schema(self, type_id: str, schema: dict[str, Any]) -> None:
228236

229237
def get(self, entity_id: str) -> GtsEntity | None:
230238
"""
231-
Get a JsonEntity by its ID.
232-
If not found in cache, try to fetch from reader.
233-
Returns None if not found.
239+
Get a JsonEntity by its ID for INTERNAL use (validation, $ref/existence
240+
resolution). The staging overlay is consulted first so an entity being
241+
validated as part of a batch can resolve its not-yet-committed siblings
242+
regardless of order. Public/API reads MUST use :meth:`get_committed` so
243+
uncommitted entities are never exposed.
234244
235245
Lookups are normalized to the canonical bare form here, so callers may
236246
pass either a bare ``gts.`` id or a ``gts://`` URI without stripping the
237247
scheme themselves.
238248
"""
239249
entity_id = strip_scheme(entity_id)
250+
with self._lock:
251+
staged = self._staged.get(entity_id)
252+
if staged is not None:
253+
return copy.deepcopy(staged)
254+
return self.get_committed(entity_id)
255+
256+
def get_committed(self, entity_id: str) -> GtsEntity | None:
257+
"""
258+
Committed-only lookup: the staging overlay is never consulted, so a
259+
staged-but-not-yet-committed entity is invisible here. This is the read
260+
path for public/API consumers. Falls back to the reader for ids not yet
261+
cached. Returns None if not found.
262+
"""
263+
entity_id = strip_scheme(entity_id)
240264
with self._lock:
241265
entity = self._by_id.get(entity_id)
242266
if entity is not None:
@@ -253,6 +277,47 @@ def get(self, entity_id: str) -> GtsEntity | None:
253277

254278
return None
255279

280+
@staticmethod
281+
def _entity_key(entity: GtsEntity) -> str:
282+
"""The registry key for an entity: the raw id for instances that carry
283+
one (anonymous/UUID instances), otherwise the canonical GTS id."""
284+
if not entity.is_schema and entity.raw_id:
285+
return entity.raw_id
286+
if entity.gts_id and entity.gts_id.id:
287+
return entity.gts_id.id
288+
if entity.raw_id:
289+
return entity.raw_id
290+
raise ValueError("Entity must have a valid gts_id or raw_id")
291+
292+
def stage(self, entity: GtsEntity) -> str:
293+
"""Place an entity into the staging overlay WITHOUT publishing it, and
294+
return its registry key. A staged entity is visible to internal
295+
validation but invisible to public reads until :meth:`commit`. Raises
296+
``EntityConflictError`` when a committed entity with different content
297+
already holds the key and updates are not allowed. Callers MUST
298+
eventually :meth:`commit` or :meth:`discard` the staged key."""
299+
stored = copy.deepcopy(entity)
300+
key = self._entity_key(stored)
301+
with self._lock:
302+
self._staged[key] = stored
303+
self._invalidate_reference_registry()
304+
return key
305+
306+
def commit(self, entity_id: str) -> None:
307+
"""Publish a previously staged entity, making it visible to public reads."""
308+
with self._lock:
309+
staged = self._staged.pop(entity_id, None)
310+
if staged is not None:
311+
self._by_id[entity_id] = staged
312+
self._invalidate_reference_registry()
313+
314+
def discard(self, entity_id: str) -> None:
315+
"""Drop a staged entity that failed validation. The committed state is
316+
untouched, so a client never observes the discarded (invalid) entity."""
317+
with self._lock:
318+
if self._staged.pop(entity_id, None) is not None:
319+
self._invalidate_reference_registry()
320+
256321
def get_schema_content(self, type_id: str) -> dict[str, Any]:
257322
"""Get schema content as dict (legacy method for backward compatibility)."""
258323
entity = self.get(type_id)
@@ -265,7 +330,11 @@ def _create_reference_registry(self) -> Registry:
265330
if self._reference_registry is not None:
266331
return self._reference_registry
267332
registry = Registry()
268-
for entity_id, entity in self._by_id.items():
333+
# Committed entities plus the staging overlay (staged overrides a
334+
# committed entry of the same id), so a schema being validated as
335+
# part of a batch resolves `$ref`s to its not-yet-committed siblings.
336+
merged = {**self._by_id, **self._staged}
337+
for entity_id, entity in merged.items():
269338
if entity.is_schema and isinstance(entity.content, dict):
270339
resource = Resource.from_contents(
271340
_without_x_gts_ref(entity.content),

0 commit comments

Comments
 (0)