Skip to content

Commit e777ea8

Browse files
committed
feat: implement dependency injection for lazy loading of submodules and enhance module resolution
1 parent 0355be3 commit e777ea8

7 files changed

Lines changed: 294 additions & 44 deletions

File tree

‎ipfs_kit_py/__init__.py‎

Lines changed: 56 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -979,8 +979,28 @@ def get_p2p_workflow_tools():
979979

980980
# Submodule lazy imports
981981
# These are optional dependencies available as submodules in the root directory
982-
def get_ipfs_datasets():
983-
"""Lazy import of ipfs_datasets_py submodule."""
982+
def get_ipfs_datasets(*, deps: object | None = None, module_override=None):
983+
"""Lazy import of ipfs_datasets_py submodule.
984+
985+
Supports dependency injection via ``module_override`` or a deps container
986+
that implements ``get_cached``/``set_cached`` (e.g., ipfs_datasets_py.RouterDeps).
987+
"""
988+
if module_override is not None:
989+
setter = getattr(deps, "set_cached", None)
990+
if callable(setter):
991+
try:
992+
setter("ipfs_kit_py::ipfs_datasets_py", module_override)
993+
except Exception:
994+
pass
995+
return module_override
996+
getter = getattr(deps, "get_cached", None)
997+
if callable(getter):
998+
try:
999+
cached = getter("ipfs_kit_py::ipfs_datasets_py")
1000+
if cached is not None:
1001+
return cached
1002+
except Exception:
1003+
pass
9841004
try:
9851005
import sys
9861006
from pathlib import Path
@@ -991,14 +1011,40 @@ def get_ipfs_datasets():
9911011
sys.path.insert(0, str(datasets_path))
9921012

9931013
import ipfs_datasets_py
1014+
setter = getattr(deps, "set_cached", None)
1015+
if callable(setter):
1016+
try:
1017+
setter("ipfs_kit_py::ipfs_datasets_py", ipfs_datasets_py)
1018+
except Exception:
1019+
pass
9941020
return ipfs_datasets_py
9951021
except ImportError as e:
9961022
logger.debug(f"ipfs_datasets_py submodule not available: {e}")
9971023
return None
9981024

9991025

1000-
def get_ipfs_accelerate():
1001-
"""Lazy import of ipfs_accelerate_py submodule."""
1026+
def get_ipfs_accelerate(*, deps: object | None = None, module_override=None):
1027+
"""Lazy import of ipfs_accelerate_py submodule.
1028+
1029+
Supports dependency injection via ``module_override`` or a deps container
1030+
that implements ``get_cached``/``set_cached`` (e.g., ipfs_datasets_py.RouterDeps).
1031+
"""
1032+
if module_override is not None:
1033+
setter = getattr(deps, "set_cached", None)
1034+
if callable(setter):
1035+
try:
1036+
setter("ipfs_kit_py::ipfs_accelerate_py", module_override)
1037+
except Exception:
1038+
pass
1039+
return module_override
1040+
getter = getattr(deps, "get_cached", None)
1041+
if callable(getter):
1042+
try:
1043+
cached = getter("ipfs_kit_py::ipfs_accelerate_py")
1044+
if cached is not None:
1045+
return cached
1046+
except Exception:
1047+
pass
10021048
try:
10031049
import sys
10041050
from pathlib import Path
@@ -1009,6 +1055,12 @@ def get_ipfs_accelerate():
10091055
sys.path.insert(0, str(accelerate_path))
10101056

10111057
import ipfs_accelerate_py
1058+
setter = getattr(deps, "set_cached", None)
1059+
if callable(setter):
1060+
try:
1061+
setter("ipfs_kit_py::ipfs_accelerate_py", ipfs_accelerate_py)
1062+
except Exception:
1063+
pass
10121064
return ipfs_accelerate_py
10131065
except ImportError as e:
10141066
logger.debug(f"ipfs_accelerate_py submodule not available: {e}")

‎ipfs_kit_py/arrow_metadata_index.py‎

Lines changed: 41 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -64,21 +64,34 @@
6464
get_ipfs_datasets_manager = None
6565
logger.info("ipfs_datasets_py not available - dataset storage disabled")
6666

67-
# Import ipfs_accelerate_py for compute acceleration
68-
try:
69-
import sys
70-
from pathlib import Path as PathlibPath
71-
accelerate_path = PathlibPath(__file__).parent.parent / "ipfs_accelerate_py"
72-
if accelerate_path.exists():
73-
sys.path.insert(0, str(accelerate_path))
74-
75-
from ipfs_accelerate_py import AccelerateCompute
76-
HAS_ACCELERATE = True
77-
logger.info("ipfs_accelerate_py compute layer available for arrow metadata")
78-
except ImportError:
79-
HAS_ACCELERATE = False
80-
AccelerateCompute = None
81-
logger.info("ipfs_accelerate_py not available - using default compute")
67+
HAS_ACCELERATE = False
68+
AccelerateCompute = None
69+
70+
71+
def _load_accelerate_compute_class(*, deps: object | None = None):
72+
"""Best-effort lazy loader for AccelerateCompute.
73+
74+
Supports dependency injection by allowing callers to pass a deps container
75+
that may already cache an imported ipfs_accelerate_py module.
76+
"""
77+
78+
global HAS_ACCELERATE, AccelerateCompute
79+
if AccelerateCompute is not None:
80+
return AccelerateCompute
81+
try:
82+
from ipfs_kit_py import get_ipfs_accelerate
83+
84+
mod = get_ipfs_accelerate(deps=deps)
85+
if mod is None:
86+
return None
87+
cls = getattr(mod, "AccelerateCompute", None)
88+
if cls is None:
89+
return None
90+
AccelerateCompute = cls
91+
HAS_ACCELERATE = True
92+
return AccelerateCompute
93+
except Exception:
94+
return None
8295

8396
#
8497
logger = logging.getLogger(__name__)
@@ -110,7 +123,10 @@ def __init__(
110123
cluster_id: str = "default",
111124
enable_dataset_storage: bool = False,
112125
enable_compute_layer: bool = False,
113-
dataset_batch_size: int = 100
126+
dataset_batch_size: int = 100,
127+
*,
128+
deps: object | None = None,
129+
accelerate_compute=None,
114130
): # Added node_id and cluster_id
115131
# Check if PyArrow is available before proceeding
116132
if not ARROW_AVAILABLE:
@@ -152,8 +168,8 @@ def __init__(
152168
self.dataset_manager = None
153169
self._metadata_buffer = []
154170

155-
# Compute layer configuration
156-
self.enable_compute_layer = enable_compute_layer and HAS_ACCELERATE
171+
# Compute layer configuration
172+
self.enable_compute_layer = bool(enable_compute_layer)
157173
self.compute_layer = None
158174

159175
# Initialize dataset manager if enabled
@@ -168,7 +184,13 @@ def __init__(
168184
# Initialize compute layer if enabled
169185
if self.enable_compute_layer:
170186
try:
171-
self.compute_layer = AccelerateCompute()
187+
if accelerate_compute is not None:
188+
self.compute_layer = accelerate_compute
189+
else:
190+
cls = _load_accelerate_compute_class(deps=deps)
191+
if cls is None:
192+
raise ImportError("ipfs_accelerate_py not available")
193+
self.compute_layer = cls()
172194
logger.info("Arrow Metadata Index compute layer enabled")
173195
except Exception as e:
174196
logger.warning(f"Failed to initialize compute layer: {e}")

‎ipfs_kit_py/deps_resolver.py‎

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,67 @@
1+
from __future__ import annotations
2+
3+
import importlib
4+
from typing import Any
5+
6+
7+
def deps_get(deps: object | None, key: str) -> Any | None:
8+
if deps is None or not key:
9+
return None
10+
getter = getattr(deps, "get_cached", None)
11+
if callable(getter):
12+
try:
13+
return getter(key)
14+
except Exception:
15+
return None
16+
if isinstance(deps, dict):
17+
return deps.get(key)
18+
return None
19+
20+
21+
def deps_set(deps: object | None, key: str, value: Any) -> Any:
22+
if deps is None or not key:
23+
return value
24+
setter = getattr(deps, "set_cached", None)
25+
if callable(setter):
26+
try:
27+
return setter(key, value)
28+
except Exception:
29+
return value
30+
if isinstance(deps, dict):
31+
deps[key] = value
32+
return value
33+
34+
35+
def resolve_module(
36+
module_name: str,
37+
*,
38+
deps: object | None = None,
39+
module_override: Any | None = None,
40+
cache_key: str | None = None,
41+
) -> Any | None:
42+
"""Resolve a Python module with optional injection + deps caching."""
43+
if module_override is not None:
44+
deps_set(deps, cache_key or f"pip::{module_name}", module_override)
45+
return module_override
46+
47+
cached = deps_get(deps, cache_key or f"pip::{module_name}")
48+
if cached is not None:
49+
return cached
50+
51+
try:
52+
module = importlib.import_module(module_name)
53+
except Exception:
54+
return None
55+
56+
deps_set(deps, cache_key or f"pip::{module_name}", module)
57+
return module
58+
59+
60+
def once(deps: object | None, key: str) -> bool:
61+
"""Return True the first time per deps container, False thereafter."""
62+
if deps is None or not key:
63+
return True
64+
if deps_get(deps, key):
65+
return False
66+
deps_set(deps, key, True)
67+
return True

‎ipfs_kit_py/enhanced_bucket_index.py‎

Lines changed: 34 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -37,20 +37,28 @@
3737
get_ipfs_datasets_manager = None
3838
logger.info("ipfs_datasets_py not available - dataset storage disabled")
3939

40-
# Import ipfs_accelerate_py for compute acceleration
41-
try:
42-
import sys
43-
accelerate_path = Path(__file__).parent.parent / "ipfs_accelerate_py"
44-
if accelerate_path.exists():
45-
sys.path.insert(0, str(accelerate_path))
46-
47-
from ipfs_accelerate_py import AccelerateCompute
48-
HAS_ACCELERATE = True
49-
logger.info("ipfs_accelerate_py compute layer available")
50-
except ImportError:
51-
HAS_ACCELERATE = False
52-
AccelerateCompute = None
53-
logger.info("ipfs_accelerate_py not available - using default compute")
40+
HAS_ACCELERATE = False
41+
AccelerateCompute = None
42+
43+
44+
def _load_accelerate_compute_class(*, deps: object | None = None):
45+
global HAS_ACCELERATE, AccelerateCompute
46+
if AccelerateCompute is not None:
47+
return AccelerateCompute
48+
try:
49+
from ipfs_kit_py import get_ipfs_accelerate
50+
51+
mod = get_ipfs_accelerate(deps=deps)
52+
if mod is None:
53+
return None
54+
cls = getattr(mod, "AccelerateCompute", None)
55+
if cls is None:
56+
return None
57+
AccelerateCompute = cls
58+
HAS_ACCELERATE = True
59+
return AccelerateCompute
60+
except Exception:
61+
return None
5462

5563
@dataclass
5664
class BucketMetadata:
@@ -93,7 +101,10 @@ def __init__(
93101
bucket_vfs_manager=None,
94102
enable_dataset_storage: bool = False,
95103
enable_compute_layer: bool = False,
96-
dataset_batch_size: int = 100
104+
dataset_batch_size: int = 100,
105+
*,
106+
deps: object | None = None,
107+
accelerate_compute=None,
97108
):
98109
"""
99110
Initialize the enhanced bucket index.
@@ -123,7 +134,7 @@ def __init__(
123134
self._index_buffer = []
124135

125136
# Compute layer configuration
126-
self.enable_compute_layer = enable_compute_layer and HAS_ACCELERATE
137+
self.enable_compute_layer = bool(enable_compute_layer)
127138
self.compute_layer = None
128139

129140
# Initialize dataset manager if enabled
@@ -138,7 +149,13 @@ def __init__(
138149
# Initialize compute layer if enabled
139150
if self.enable_compute_layer:
140151
try:
141-
self.compute_layer = AccelerateCompute()
152+
if accelerate_compute is not None:
153+
self.compute_layer = accelerate_compute
154+
else:
155+
cls = _load_accelerate_compute_class(deps=deps)
156+
if cls is None:
157+
raise ImportError("ipfs_accelerate_py not available")
158+
self.compute_layer = cls()
142159
logger.info("Enhanced Bucket Index compute layer enabled")
143160
except Exception as e:
144161
logger.warning(f"Failed to initialize compute layer: {e}")

‎ipfs_kit_py/ipfs_kit.py‎

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -372,6 +372,10 @@ def __init__(
372372
enable_cluster_management=False,
373373
enable_metadata_index=False,
374374
auto_start_daemons=True,
375+
*,
376+
deps: object | None = None,
377+
ipfs_accelerate=None,
378+
ipfs_datasets=None,
375379
):
376380
"""
377381
Initializes the IPFS Kit instance.
@@ -408,6 +412,24 @@ def __init__(
408412
# passing metadata.
409413
metadata = self.metadata
410414

415+
# Dependency injection container and optional injected cross-package modules.
416+
self.deps = deps or (resources.get("deps") if isinstance(resources, dict) else None) or self.metadata.get("deps")
417+
self.ipfs_accelerate_py = ipfs_accelerate or (
418+
resources.get("ipfs_accelerate_py") if isinstance(resources, dict) else None
419+
)
420+
self.ipfs_datasets_py = ipfs_datasets or (resources.get("ipfs_datasets_py") if isinstance(resources, dict) else None)
421+
422+
# Best-effort: cache injected modules on deps for other packages to reuse.
423+
setter = getattr(self.deps, "set_cached", None)
424+
if callable(setter):
425+
try:
426+
if self.ipfs_accelerate_py is not None:
427+
setter("ipfs_kit_py::ipfs_accelerate_py", self.ipfs_accelerate_py)
428+
if self.ipfs_datasets_py is not None:
429+
setter("ipfs_kit_py::ipfs_datasets_py", self.ipfs_datasets_py)
430+
except Exception:
431+
pass
432+
411433
# Initialize _initialized attribute
412434
self._initialized = False
413435

@@ -529,9 +551,8 @@ def __init__(
529551
self.huggingface_kit = huggingface_kit(resources=resources, metadata=metadata)
530552
else:
531553
self.huggingface_kit = None # Initialize to None
532-
554+
533555
# Initialize ipget component
534-
self.ipget = ipget(resources=resources, metadata={"role": self.role})
535556

536557
# Initialize Lotus Kit if available and not disabled
537558
if HAS_LOTUS and "lotus" not in disabled_components:

0 commit comments

Comments
 (0)