feat(file-summary): 实现文件处理技能链路
This commit is contained in:
@@ -16,19 +16,34 @@ from review_agent.models import (
|
||||
)
|
||||
|
||||
from .events import record_event
|
||||
from .skills.archive_extract import ArchiveExtractSkill
|
||||
from .skills.base import WorkflowContext
|
||||
from .skills.document_page_count import DocumentPageCountSkill
|
||||
from .skills.file_inventory import FileInventorySkill
|
||||
from .skills.product_detect import ProductDetectSkill
|
||||
from .skills.registry import SkillRegistry
|
||||
|
||||
|
||||
NODE_DEFINITIONS = [
|
||||
("upload", "附件固化"),
|
||||
("extract", "压缩包解压"),
|
||||
("inventory", "文件扫描"),
|
||||
("page_count", "页数统计"),
|
||||
("product_detect", "产品识别"),
|
||||
("report", "报告输出"),
|
||||
("complete", "完成"),
|
||||
("upload", "附件固化", ""),
|
||||
("extract", "压缩包解压", "archive_extract"),
|
||||
("inventory", "文件扫描", "file_inventory"),
|
||||
("page_count", "页数统计", "document_page_count"),
|
||||
("product_detect", "产品识别", "product_detect"),
|
||||
("report", "报告输出", ""),
|
||||
("complete", "完成", ""),
|
||||
]
|
||||
|
||||
|
||||
def default_skill_registry() -> SkillRegistry:
|
||||
registry = SkillRegistry()
|
||||
registry.register(ArchiveExtractSkill())
|
||||
registry.register(FileInventorySkill())
|
||||
registry.register(DocumentPageCountSkill())
|
||||
registry.register(ProductDetectSkill())
|
||||
return registry
|
||||
|
||||
|
||||
def build_batch_no() -> str:
|
||||
return f"FS-{timezone.localtime().strftime('%Y%m%d%H%M%S')}-{uuid4().hex[:6]}"
|
||||
|
||||
@@ -61,7 +76,7 @@ def create_file_summary_batch(
|
||||
attachment.upload_status = FileAttachment.UploadStatus.BOUND
|
||||
attachment.save(update_fields=["upload_status"])
|
||||
|
||||
for code, name in NODE_DEFINITIONS:
|
||||
for code, name, _skill_name in NODE_DEFINITIONS:
|
||||
WorkflowNodeRun.objects.create(batch=batch, node_code=code, node_name=name)
|
||||
|
||||
record_event(batch, "workflow_created", {"batch_id": batch.pk, "batch_no": batch.batch_no})
|
||||
@@ -69,8 +84,9 @@ def create_file_summary_batch(
|
||||
|
||||
|
||||
class WorkflowExecutor:
|
||||
def __init__(self, batch: FileSummaryBatch):
|
||||
def __init__(self, batch: FileSummaryBatch, registry: SkillRegistry | None = None):
|
||||
self.batch = batch
|
||||
self.registry = registry or default_skill_registry()
|
||||
|
||||
def run(self) -> None:
|
||||
self.batch.status = FileSummaryBatch.Status.RUNNING
|
||||
@@ -107,6 +123,15 @@ class WorkflowExecutor:
|
||||
{"node_code": node.node_code, "status": node.status, "progress": node.progress},
|
||||
)
|
||||
|
||||
skill_name = next(
|
||||
(skill for code, _name, skill in NODE_DEFINITIONS if code == node.node_code),
|
||||
"",
|
||||
)
|
||||
if skill_name:
|
||||
result = self.registry.execute(skill_name, WorkflowContext(batch=self.batch))
|
||||
if not result.success:
|
||||
raise RuntimeError(result.message or f"{node.node_name}执行失败")
|
||||
|
||||
node.status = WorkflowNodeRun.Status.SUCCESS
|
||||
node.progress = 100
|
||||
node.finished_at = timezone.now()
|
||||
|
||||
Reference in New Issue
Block a user