This commit is contained in:
jinli.yl 2025-07-09 20:52:05 +08:00
parent 0bb3068ad2
commit d82af71fed
13 changed files with 39 additions and 38 deletions

View file

@ -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" \

View file

@ -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)

View file

@ -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"

View file

@ -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]):

View file

@ -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

View file

@ -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):

View file

@ -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

View file

@ -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()

View file

@ -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

View file

@ -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)

View file

@ -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:

View file

@ -1,4 +1,3 @@
from abc import ABC
from typing import List
from pydantic import BaseModel, Field

View file

@ -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):