From d82af71fedf01e71538a5d5f309c14bd6f2e6a77 Mon Sep 17 00:00:00 2001 From: "jinli.yl" Date: Wed, 9 Jul 2025 20:52:05 +0800 Subject: [PATCH] up --- cookbook/vector_store/elasticsearch.md | 3 ++- v1/app.py | 6 +++-- v1/config/demo_config.yaml | 28 ++++++++++++---------- v1/embedding_model/base_embedding_model.py | 8 ------- v1/op/__init__.py | 2 ++ v1/op/base_op.py | 2 +- v1/op/mock_op.py | 2 +- v1/pipeline/pipeline.py | 8 +++---- v1/pipeline/pipeline_context.py | 6 ++--- v1/schema/app_config.py | 2 +- v1/schema/message.py | 3 ++- v1/schema/request.py | 1 - v1/utils/timer.py | 6 ++--- 13 files changed, 39 insertions(+), 38 deletions(-) diff --git a/cookbook/vector_store/elasticsearch.md b/cookbook/vector_store/elasticsearch.md index 6b2b71ba..3df0b3b0 100644 --- a/cookbook/vector_store/elasticsearch.md +++ b/cookbook/vector_store/elasticsearch.md @@ -28,7 +28,8 @@ docker run -p 9200:9200 \ #### Docker Run Image with Http Host ```shell docker pull docker.elastic.co/elasticsearch/elasticsearch-wolfi:9.0.0 -docker run -p 8200:9200 \ +docker run -p 9201:9201 \ + --memory='4GB' \ -e "discovery.type=single-node" \ -e "xpack.security.enabled=false" \ -e "xpack.license.self_generated.type=trial" \ diff --git a/v1/app.py b/v1/app.py index e88de5e6..a89a1c32 100644 --- a/v1/app.py +++ b/v1/app.py @@ -1,12 +1,15 @@ import sys import uvicorn +from dotenv import load_dotenv from fastapi import FastAPI from v1.schema.request import RetrieverRequest, SummarizerRequest, VectorStoreRequest, AgentRequest from v1.schema.response import RetrieverResponse, SummarizerResponse, VectorStoreResponse, AgentResponse from v1.service.experience_maker_service import ExperienceMakerService +load_dotenv() + app = FastAPI() service = ExperienceMakerService(sys.argv[1:]) @@ -35,5 +38,4 @@ if __name__ == "__main__": host=service.http_service_config.host, port=service.http_service_config.port, timeout_keep_alive=service.http_service_config.timeout_keep_alive, - limit_concurrency=service.http_service_config.limit_concurrency, - workers=service.http_service_config.workers) + limit_concurrency=service.http_service_config.limit_concurrency) diff --git a/v1/config/demo_config.yaml b/v1/config/demo_config.yaml index af3aa8c8..9852af95 100644 --- a/v1/config/demo_config.yaml +++ b/v1/config/demo_config.yaml @@ -5,26 +5,27 @@ http_service: port: 8001 timeout_keep_alive: 600 limit_concurrency: 64 - workers: 2 thread_pool: - max_workers: 20 + max_workers: 10 api: - step_retriever: mock1_op->mock2_op->mock3_op - step_summarizer: mock1_op->[mock4_op->mock2_op|mock5_op]->mock3_op + retriever: mock1_op->[mock4_op->mock2_op|mock5_op]->[mock3_op|mock6_op] + summarizer: mock1_op->[mock4_op->mock2_op|mock5_op]->mock3_op vector_store: mock6_op op: mock1_op: backend: mock1_op - a: 1 - b: 2 llm: default vector_store: default + params: + a: 1 + b: 2 mock2_op: backend: mock2_op - a: 1 + params: + a: 1 mock3_op: backend: mock3_op mock4_op: @@ -37,18 +38,21 @@ op: llm: default: backend: openai_compatible - model_name: qwen3-32b - temperature: 0.6 + model: qwen3-32b + params: + temperature: 0.6 embedding_model: default: backend: openai_compatible - model_name: text-embedding-v4 - dimensions: 1024 + model: text-embedding-v4 + params: + dimensions: 1024 vector_store: default: backend: elasticsearch embedding_model: default - hosts: "http://localhost:9200" + params: + hosts: "http://localhost:9200" diff --git a/v1/embedding_model/base_embedding_model.py b/v1/embedding_model/base_embedding_model.py index 0c80405f..a6a96fd6 100644 --- a/v1/embedding_model/base_embedding_model.py +++ b/v1/embedding_model/base_embedding_model.py @@ -14,14 +14,6 @@ class BaseEmbeddingModel(BaseModel, ABC): raise_exception: bool = Field(default=True, description="raise exception") def _get_embeddings(self, input_text: str | List[str]): - """ - Get the embedding vector based on the input text. - This is an abstract method, and its concrete implementation must be provided in a subclass to generate the embedding vector for the given text. - Args: - input_text (str | List[str]): The input text, which can be a single string or a list of strings. - Raises: - NotImplementedError: If the method is not implemented in the subclass. - """ raise NotImplementedError def get_embeddings(self, input_text: str | List[str]): diff --git a/v1/op/__init__.py b/v1/op/__init__.py index b8c42b57..d0ec5274 100644 --- a/v1/op/__init__.py +++ b/v1/op/__init__.py @@ -1,3 +1,5 @@ from v1.utils.registry import Registry OP_REGISTRY = Registry() + +from v1.op.mock_op import MockOp1, MockOp2, MockOp3, MockOp4, MockOp5, MockOp6 diff --git a/v1/op/base_op.py b/v1/op/base_op.py index ce08a31a..ca6afc27 100644 --- a/v1/op/base_op.py +++ b/v1/op/base_op.py @@ -15,7 +15,7 @@ class BaseOp(ABC): @property def simple_name(self) -> str: - return self.__class__.__name__.lower().replace("op", "") + return self.__class__.__name__.lower() @abstractmethod def execute(self): diff --git a/v1/op/mock_op.py b/v1/op/mock_op.py index 7449011c..96e01628 100644 --- a/v1/op/mock_op.py +++ b/v1/op/mock_op.py @@ -9,7 +9,7 @@ from v1.op.base_op import BaseOp @OP_REGISTRY.register("mock1_op") class MockOp1(BaseOp): - def __init__(self, a: int, b: str, **kwargs): + def __init__(self, a: int = 1, b: str = "2", **kwargs): super().__init__(**kwargs) self.a = a self.b = b diff --git a/v1/pipeline/pipeline.py b/v1/pipeline/pipeline.py index a2853229..eb7dc1ea 100644 --- a/v1/pipeline/pipeline.py +++ b/v1/pipeline/pipeline.py @@ -28,10 +28,10 @@ class Pipeline: continue if self.parallel_symbol in sub_pipeline: - self.pipeline_list.append(sub_pipeline.split(self.parallel_symbol)) + pipeline_list.append(sub_pipeline.split(self.parallel_symbol)) else: - self.pipeline_list.append(sub_pipeline) - + pipeline_list.append(sub_pipeline) + logger.info(f"add sub_pipeline={sub_pipeline}") return pipeline_list def _execute_sub_pipeline(self, pipeline: str): @@ -74,7 +74,7 @@ class Pipeline: else: raise ValueError(f"unknown pipeline.type={type(pipeline)}") - @timer() + @timer(name="pipeline.execute") def __call__(self, enable_print: bool = True): if enable_print: self.print_pipeline() diff --git a/v1/pipeline/pipeline_context.py b/v1/pipeline/pipeline_context.py index a3dc597e..ccdd92ed 100644 --- a/v1/pipeline/pipeline_context.py +++ b/v1/pipeline/pipeline_context.py @@ -5,15 +5,15 @@ from v1.schema.app_config import AppConfig from v1.vector_store.base_vector_store import BaseVectorStore -class PipelineContext(object): +class PipelineContext: def __init__(self, **kwargs): self._context: dict = {**kwargs} - def __getattr__(self, key: str, default=None): + def get_context(self, key: str, default=None): return self._context.get(key, default) - def __setattr__(self, key: str, value): + def set_context(self, key: str, value): self._context[key] = value @property diff --git a/v1/schema/app_config.py b/v1/schema/app_config.py index dfcfe876..d5c63acd 100644 --- a/v1/schema/app_config.py +++ b/v1/schema/app_config.py @@ -8,7 +8,6 @@ class HttpServiceConfig: port: int = field(default=8001) timeout_keep_alive: int = field(default=600) limit_concurrency: int = field(default=64) - workers: int | None = field(default=None) @dataclass @@ -54,6 +53,7 @@ class VectorStoreConfig: params: dict = field(default_factory=dict) +@dataclass class AppConfig: config_path: str = field(default="") http_service: HttpServiceConfig = field(default_factory=HttpServiceConfig) diff --git a/v1/schema/message.py b/v1/schema/message.py index 9f431cc6..786fed14 100644 --- a/v1/schema/message.py +++ b/v1/schema/message.py @@ -1,7 +1,7 @@ import json from typing import List -from pydantic import BaseModel, Field +from pydantic import BaseModel, Field, model_validator from v1.enumeration.role import Role @@ -48,6 +48,7 @@ class Message(BaseModel): metadata: dict = Field(default_factory=dict) def simple_dump(self, add_reason_when_empty: bool = True) -> dict: + result: dict if self.content: result = {"role": self.role.value, "content": self.content} elif add_reason_when_empty and self.reasoning_content: diff --git a/v1/schema/request.py b/v1/schema/request.py index 42206d32..b99f476c 100644 --- a/v1/schema/request.py +++ b/v1/schema/request.py @@ -1,4 +1,3 @@ -from abc import ABC from typing import List from pydantic import BaseModel, Field diff --git a/v1/utils/timer.py b/v1/utils/timer.py index 318b59cf..5daa87ce 100644 --- a/v1/utils/timer.py +++ b/v1/utils/timer.py @@ -16,7 +16,7 @@ class Timer(object): def __enter__(self, *args, **kwargs): self.time_start = time.time() - logger.info(f"========== timer={self.name} start ==========", stacklevel=self.stack_level) + logger.info(f"========== {self.name} start ==========", stacklevel=self.stack_level) return self def __exit__(self, *args): @@ -27,10 +27,10 @@ class Timer(object): else: time_str = f"{self.time_cost:.3f}s" - logger.info(f"========== timer={self.name} end, time_cost={time_str} ==========", stacklevel=self.stack_level) + logger.info(f"========== {self.name} end, time_cost={time_str} ==========", stacklevel=self.stack_level) -def timer(name: Optional[str] = None, use_ms: bool = False, stack_level: int = 2): +def timer(name: str = None, use_ms: bool = False, stack_level: int = 2): def decorator(func): def wrapper(*args, **kwargs): with Timer(name=name or func.__name__, use_ms=use_ms, stack_level=stack_level + 1):