Skip to content

Commit ec5aa1e

Browse files
refactor: Fix MRO of schema stream attribute
Signed-off-by: Edgar Ramírez Mondragón <edgarrm358@gmail.com>
1 parent 8a4f294 commit ec5aa1e

4 files changed

Lines changed: 91 additions & 44 deletions

File tree

‎singer_sdk/contrib/filesystem/stream.py‎

Lines changed: 1 addition & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,6 @@
33
from __future__ import annotations
44

55
import abc
6-
import functools
76
import sys
87
import typing as t
98

@@ -66,6 +65,7 @@ def __init__(
6665
self.filesystem = filesystem
6766

6867
super().__init__(tap, schema=None, name=name)
68+
self._schema = self._get_full_schema()
6969

7070
# TODO(edgarrmondragon): Make this None if the filesystem does not support it.
7171
self.replication_key = SDC_META_MODIFIED_AT
@@ -89,12 +89,6 @@ def _get_full_schema(self) -> dict[str, t.Any]:
8989
schema["properties"].update(self.SDC_PROPERTIES)
9090
return schema
9191

92-
@functools.cached_property
93-
@override
94-
def schema(self) -> dict[str, t.Any]:
95-
"""Return the schema for the stream."""
96-
return self._get_full_schema()
97-
9892
@override
9993
def get_records(
10094
self,

‎singer_sdk/schema/source.py‎

Lines changed: 78 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,11 @@
2929
else:
3030
from typing_extensions import assert_never
3131

32+
if sys.version_info >= (3, 11):
33+
from typing import Self # noqa: ICN003
34+
else:
35+
from typing_extensions import Self
36+
3237
if sys.version_info >= (3, 12):
3338
from typing import override # noqa: ICN003
3439
else:
@@ -74,6 +79,7 @@ def preprocess_schema(
7479
Returns:
7580
The preprocessed schema.
7681
"""
82+
...
7783

7884

7985
class SchemaSource(ABC, t.Generic[_TKey]):
@@ -186,7 +192,7 @@ class MyStream(Stream):
186192

187193
def __init__(
188194
self,
189-
schema_source: SchemaSource[_TKey],
195+
schema_source: SchemaSource[_TKey] | None = None,
190196
*,
191197
key: _TKey | None = None,
192198
) -> None:
@@ -200,18 +206,37 @@ def __init__(
200206
self.schema_source = schema_source
201207
self.key = key
202208

209+
@t.overload
210+
def __get__(self, obj: None, objtype: type[Stream] | None = ...) -> Self: ...
211+
@t.overload
212+
def __get__(self, obj: Stream, objtype: type[Stream] | None = ...) -> Schema: ...
203213
@t.final
204-
def __get__(self, obj: Stream, objtype: type[Stream]) -> Schema:
214+
def __get__(
215+
self,
216+
obj: Stream | None,
217+
objtype: type[Stream] | None = None,
218+
) -> Schema | Self:
205219
"""Get the schema from the schema source.
206220
207221
Args:
208-
obj: The object to get the schema from.
209-
objtype: The type of the object to get the schema from.
222+
obj: The stream instance, or ``None`` when accessed from the class.
223+
objtype: The stream class.
210224
211225
Returns:
212-
A JSON schema dictionary.
226+
A JSON schema dictionary, or the descriptor itself for class-level access.
227+
"""
228+
if obj is None:
229+
return self
230+
return self.get_stream_schema(obj, objtype or type(obj))
231+
232+
def __set__(self, obj: Stream, _value: object) -> None:
233+
"""Raise AttributeError; schema is read-only on instances.
234+
235+
Raises:
236+
AttributeError: Always.
213237
"""
214-
return self.get_stream_schema(obj, objtype)
238+
msg = f"'schema' is read-only on {type(obj).__name__}"
239+
raise AttributeError(msg)
215240

216241
def get_stream_schema(
217242
self,
@@ -226,13 +251,60 @@ def get_stream_schema(
226251
227252
Returns:
228253
A JSON schema dictionary.
254+
255+
Raises:
256+
DiscoveryError: If no schema source is configured and this method is not
257+
overridden.
229258
"""
259+
if self.schema_source is None:
260+
msg = (
261+
f"No schema source configured for {type(self).__name__!r}; "
262+
"either pass a schema_source or override get_stream_schema()."
263+
)
264+
raise DiscoveryError(msg)
230265
return self.schema_source.get_schema(
231266
self.key or stream.name, # type: ignore[arg-type] # ty:ignore[invalid-argument-type]
232267
key_properties=stream.primary_keys,
233268
)
234269

235270

271+
class _InstanceSchemaDescriptor(StreamSchema[str]):
272+
"""Fallback descriptor placed on the base ``Stream`` class.
273+
274+
When no class-level ``schema`` is declared on a stream subclass, this
275+
descriptor is found in the MRO and forwards the attribute access to the
276+
instance's ``_schema`` attribute (set by ``Stream.__init__`` or by the
277+
catalog-loading path).
278+
"""
279+
280+
def __init__(self) -> None:
281+
"""Initialize with no schema source."""
282+
super().__init__()
283+
284+
@override
285+
def get_stream_schema(
286+
self,
287+
stream: _TStream,
288+
stream_class: type[_TStream],
289+
) -> Schema:
290+
"""Return the schema stored on the stream instance.
291+
292+
Args:
293+
stream: The stream instance.
294+
stream_class: The stream class (unused).
295+
296+
Returns:
297+
The JSON schema dictionary.
298+
299+
Raises:
300+
DiscoveryError: If no schema has been set on the instance.
301+
"""
302+
if stream._schema is None: # noqa: SLF001
303+
msg = f"The schema for stream '{stream.name}' was not provided"
304+
raise DiscoveryError(msg)
305+
return stream._schema # noqa: SLF001
306+
307+
236308
def _load_yaml(content: bytes) -> dict[str, t.Any]:
237309
import yaml # noqa: PLC0415
238310

‎singer_sdk/sql/stream.py‎

Lines changed: 1 addition & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,6 @@
44

55
import abc
66
import typing as t
7-
from functools import cached_property
87

98
import sqlalchemy as sa
109

@@ -28,8 +27,6 @@ class SQLStream(Stream, abc.ABC):
2827
"""Base class for SQLAlchemy-based streams."""
2928

3029
connector_class = SQLConnector
31-
_cached_schema: dict | None = None
32-
3330
supports_nulls_first: bool = False
3431
"""Whether the database supports the NULLS FIRST/LAST syntax."""
3532

@@ -53,7 +50,7 @@ def __init__(
5350
self.catalog_entry = catalog_entry
5451
super().__init__(
5552
tap=tap,
56-
schema=self.schema,
53+
schema=CatalogEntry.from_dict(catalog_entry).schema,
5754
name=self.tap_stream_id,
5855
)
5956

@@ -86,17 +83,6 @@ def metadata(self) -> MetadataMapping:
8683
"""
8784
return self._singer_catalog_entry.metadata
8885

89-
@cached_property
90-
def schema(self) -> dict:
91-
"""Return metadata object (dict) as specified in the Singer spec.
92-
93-
Metadata from an input catalog will override standard metadata.
94-
95-
Returns:
96-
The schema object.
97-
"""
98-
return self._singer_catalog_entry.schema.to_dict()
99-
10086
@property
10187
def tap_stream_id(self) -> str:
10288
"""Return the unique ID used by the tap to identify this stream.

‎singer_sdk/streams/core.py‎

Lines changed: 11 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919
from singer_sdk.exceptions import (
2020
AbortedSyncFailedException,
2121
AbortedSyncPausedException,
22-
DiscoveryError,
2322
InvalidReplicationKeyException,
2423
InvalidStreamSortException,
2524
MaxRecordsLimitException,
@@ -35,6 +34,7 @@
3534
from singer_sdk.helpers._util import utc_now
3635
from singer_sdk.helpers.conform import TypeConformanceLevel
3736
from singer_sdk.mapper import RemoveRecordTransform, SameRecordTransform
37+
from singer_sdk.schema.source import _InstanceSchemaDescriptor
3838
from singer_sdk.singerlib.catalog import (
3939
REPLICATION_FULL_TABLE,
4040
REPLICATION_INCREMENTAL,
@@ -47,6 +47,7 @@
4747
from singer_sdk.helpers._batch import BaseBatchFileEncoding
4848
from singer_sdk.helpers._compat import Traversable
4949
from singer_sdk.mapper import StreamMap
50+
from singer_sdk.schema.source import StreamSchema
5051
from singer_sdk.singerlib.catalog import StreamMetadata
5152
from singer_sdk.tap_base import Tap
5253

@@ -102,6 +103,15 @@ class Stream(abc.ABC): # noqa: PLR0904
102103
selected_by_default: bool = True
103104
"""Whether this stream is selected by default in the catalog."""
104105

106+
schema: t.ClassVar[StreamSchema[t.Any]] = _InstanceSchemaDescriptor()
107+
"""JSON schema for this stream.
108+
109+
Can be set at the class level to a :class:`~singer_sdk.StreamSchema` descriptor
110+
(e.g. ``StreamSchema(SchemaDirectory("schemas"))``) or to a plain
111+
``dict[str, Any]`` JSON Schema object. When neither is provided the schema
112+
must be supplied via the *schema* constructor argument.
113+
"""
114+
105115
def __init__(
106116
self,
107117
tap: Tap,
@@ -514,21 +524,6 @@ def schema_filepath(self) -> Path | Traversable | None:
514524
"""
515525
return self._schema_filepath
516526

517-
@property
518-
def schema(self) -> dict:
519-
"""Get schema.
520-
521-
Returns:
522-
JSON Schema dictionary for this stream.
523-
524-
Raises:
525-
DiscoveryError: If the schema was not provided.
526-
"""
527-
if self._schema is None:
528-
msg = f"The schema for stream '{self.name}' was not provided"
529-
raise DiscoveryError(msg)
530-
return self._schema
531-
532527
@property
533528
def primary_keys(self) -> t.Sequence[str]:
534529
"""Get primary keys.

0 commit comments

Comments
 (0)