UFO Galaxy TaskStarLine 深度解析:用智能依赖边编织任务星座 DAG
【免费下载链接】UFOUFO³: Weaving the Digital Agent Galaxy项目地址: https://gitcode.com/GitHub_Trending/uf/UFO
本篇技术指南围绕 UFO 项目 Galaxy 框架中的核心组件TaskStarLine展开,深入讲解其在任务星座(Task Constellation)中如何以有向依赖边(edge)连接两个 TaskStar,支撑条件逻辑、仅成功执行与自定义条件求值等依赖关系。读完本文,你将掌握 TaskStarLine 的四种依赖类型、生命周期状态机、条件求值引擎、序列化方案及其与 TaskConstellation 的集成方式,并能在构建多智能体工作流时正确设计无环依赖图。
TaskStarLine 是什么:任务星座 DAG 的“边”
在 Galaxy 的任务编排体系(Constellation V2)中,TaskStar是原子执行单元(顶点),而TaskStarLine则负责连接这些顶点,形成一张有向无环图(DAG)。每条 TaskStarLine 定义了一个从任务 $t_i$ 到任务 $t_j$ 的有向依赖关系,并携带依赖类型、自然语言条件描述与可编程的条件求值器。
其形式化定义为:
$$ e_{i \rightarrow j} = (\text{from_task}_i, \text{to_task}_j, \text{type}, \text{description}) $$
任务 $t_j$ 必须等到 $t_i$ 上的条件(依据依赖类型而定)被满足后,才能开始执行。例如在“checkout → build → test → deploy”流水线中,deploy 依赖于 test,而 test 又依赖于 build,每一条依赖都是一条 TaskStarLine。
从源码看,TaskStarLine 直接实现了 IDependency 接口——该接口规定了source_task_id、target_task_id、dependency_type与is_satisfied(completed_tasks)四个抽象成员。这意味着任何遵守IDependency契约的组件都能与 TaskStarLine 互换使用,从而保证了框架的接口隔离原则与可扩展性。
核心属性与状态追踪
核心属性
| 属性 | 类型 | 说明 |
|---|---|---|
| line_id | str | 唯一标识,未提供时自动生成 UUID(源码见 task_star_line.py) |
| from_task_id | str | 前置任务(源)ID |
| to_task_id | str | 依赖任务(目标)ID |
| dependency_type | DependencyType | 依赖关系类型 |
| condition_description | str | 条件的自然语言描述(调试与可读性友好) |
| condition_evaluator | Callable | 求值条件是否满足的函数,接收前置任务结果、返回bool |
| metadata | Dict[str, Any] | 依赖的附加元数据 |
兼容性别名:source_task_id与target_task_id分别是from_task_id与to_task_id的只读别名,用于满足IDependency接口兼容(task_star_line.py)。
状态追踪属性
| 属性 | 类型 | 说明 |
|---|---|---|
| is_satisfied | bool | 依赖条件当前是否满足 |
| last_evaluation_result | bool | 最近一次条件求值的结果 |
| last_evaluation_time | datetime | 最近一次条件求值的时间戳 |
| created_at | datetime | 依赖创建时间戳 |
| updated_at | datetime | 最近修改时间戳 |
所有状态追踪属性均为只读,由 TaskStarLine 内部方法自动维护:创建时created_at与updated_at初始化为同一 UTC 时间戳,每次求值、修改或手动置位时自动刷新(task_star_line.py)。时间戳统一使用datetime.now(timezone.utc),保证跨设备协调时的时区一致性。
四种依赖类型
依赖类型由 DependencyType 枚举 定义,其序列化值为小写字符串:unconditional、conditional、success_only、completion_only。
1. 无条件依赖(UNCONDITIONAL)
任务 $t_j$总是等待 $t_i$ 完成,无论成功还是失败。
适用场景:顺序流水线阶段;任何任务完成后的资源清理;日志或通知任务。
# Task B always runs after Task A completes dep = TaskStarLine.create_unconditional( from_task_id="task_a", to_task_id="task_b", description="B runs after A regardless of outcome" )在求值引擎中,UNCONDITIONAL直接短路返回True(task_star_line.py),即只要执行到求值环节就视为满足——它的语义是“执行权交给前置任务完成这个事实本身”。
2. 仅成功依赖(SUCCESS_ONLY)
任务 $t_j$仅当$t_i$ 成功完成时才继续执行。
适用场景:构建流水线(仅当构建成功才部署);多步数据处理;条件工作流分支。
# Task B only runs if Task A succeeds dep = TaskStarLine.create_success_only( from_task_id="build_task", to_task_id="deploy_task", description="Deploy only if build succeeds" )注意:成功的判定标准是前置任务返回结果不为None(task_star_line.py)。这一约定与 TaskStar 的执行结果语义保持一致——返回None即视为失败。源码实现为result = prerequisite_result is not None,是四种类型中最直观的判定。
3. 仅完成依赖(COMPLETION_ONLY)
任务 $t_j$ 在 $t_i$完成时即继续,无论成功或失败。
适用场景:清理任务;通知任务;审计日志。
# Task B runs after Task A finishes, regardless of outcome dep = TaskStarLine( from_task_id="main_task", to_task_id="cleanup_task", dependency_type=DependencyType.COMPLETION_ONLY, condition_description="Cleanup runs regardless of main task outcome" )从语义上,COMPLETION_ONLY与UNCONDITIONAL在求值结果上等价(都返回True),区别在于设计意图与文档化表达:前者强调“无论结果如何都要执行”,通常用于收尾类任务;后者强调“排队等待前置完成”。选择类型时建议按业务语义区分,便于后续阅读与调试。
4. 条件依赖(CONDITIONAL)
任务 $t_j$ 依据用户自定义条件对 $t_i$ 结果的求值结果决定是否继续。
适用场景:错误处理分支;基于结果的路由;性能驱动的优化路径。
# Define custom condition evaluator def check_coverage_threshold(result): """Run next task only if test coverage > 80%""" if result and isinstance(result, dict): coverage = result.get("coverage_percent", 0) return coverage > 80 return False # Create conditional dependency dep = TaskStarLine.create_conditional( from_task_id="test_task", to_task_id="quality_gate_task", condition_description="Proceed if test coverage > 80%", condition_evaluator=check_coverage_threshold )注意:若为CONDITIONAL类型但未提供condition_evaluator,则默认回退为SUCCESS_ONLY行为(即检查结果是否为None)。该回退逻辑在源码中清晰可见:task_star_line.py。
依赖生命周期
依赖对象从创建到被消费,遵循如下状态机:
从实现角度可以这样理解该状态机:Created对应构造完成且_is_satisfied=False;当前置任务运行完毕,星座编排器调用evaluate_condition(prerequisite_result)进入Evaluating;求值结果写入_last_evaluation_result并同步到_is_satisfied,随后进入Satisfied或Unsatisfied。mark_satisfied()可视为人工强制从任意状态跳到Satisfied。
实战用法:创建依赖
以下代码演示了四种依赖的完整创建方式,覆盖从工厂方法到手动构造的全部路径:
from galaxy.constellation import TaskStarLine from galaxy.constellation.enums import DependencyType # 1. Unconditional dependency dep1 = TaskStarLine.create_unconditional( from_task_id="checkout_code", to_task_id="build_project", description="Build after checkout" ) # 2. Success-only dependency dep2 = TaskStarLine.create_success_only( from_task_id="build_project", to_task_id="deploy_staging", description="Deploy only if build succeeds" ) # 3. Conditional dependency with custom logic def check_test_results(result): return result.get("tests_passed", 0) == result.get("total_tests", 0) dep3 = TaskStarLine.create_conditional( from_task_id="run_tests", to_task_id="deploy_production", condition_description="Deploy to production only if all tests pass", condition_evaluator=check_test_results ) # 4. Manual construction dep4 = TaskStarLine( from_task_id="task_a", to_task_id="task_b", dependency_type=DependencyType.COMPLETION_ONLY, condition_description="Task B runs after Task A completes", metadata={"priority": "high", "category": "cleanup"} )三个工厂方法create_unconditional、create_success_only、create_conditional均为类方法(classmethod),其签名与默认描述在 task_star_line.py 中有精确定义。手动构造则提供了最大自由度,可直接传入line_id与metadata。
核心操作
条件求值
# Evaluate condition with prerequisite result prerequisite_result = { "status": "success", "coverage_percent": 85, "tests_passed": 120, "total_tests": 120 } is_satisfied = dep.evaluate_condition(prerequisite_result) if is_satisfied: print("✅ Dependency satisfied, dependent task can run") print(f"Evaluated at: {dep.last_evaluation_time}") else: print("❌ Dependency not satisfied, dependent task blocked") # Check evaluation history print(f"Last result: {dep.last_evaluation_result}")evaluate_condition是依赖引擎的心脏,其实现(task_star_line.py)依次执行以下步骤:
- 记录求值时间戳
_last_evaluation_time; - 依据
dependency_type分支求值:UNCONDITIONAL/COMPLETION_ONLY直接为True,SUCCESS_ONLY判定结果非None,CONDITIONAL调用用户求值器(无求值器则回退成功判定); - 将结果写入
_last_evaluation_result与_is_satisfied; - 异常兜底:整个求值体包裹在
try/except中,求值器抛出任何异常都会吞掉并返回False(同时记录错误但不上抛)。这一设计保证了星座执行引擎不会被用户求值器中的意外异常击穿。
手动满足控制
# Manually mark dependency as satisfied (override) dep.mark_satisfied() # Reset satisfaction status dep.reset_satisfaction() # Check satisfaction if dep.is_satisfied(): print("Dependency is satisfied")mark_satisfied()会将_is_satisfied、_last_evaluation_result置为True并刷新时间戳;reset_satisfaction()则清空全部求值状态。二者常用于调试、强制放行或回滚场景。
状态查询
is_satisfied()具有双模式语义,理解这一点至关重要:
# Method 1: Check using completed tasks list (for IDependency interface) # Returns True if from_task_id is in the completed_tasks list completed_tasks = ["task_a", "task_b", "task_c"] if dep.is_satisfied(completed_tasks): print("Prerequisite task is completed") # Method 2: Check internal satisfaction state (without parameter) # Returns the internal _is_satisfied flag set by evaluate_condition if dep.is_satisfied(): print("Dependency condition is satisfied") # Get last evaluation details print(f"Last evaluated: {dep.last_evaluation_time}") print(f"Result: {dep.last_evaluation_result}") # Access metadata print(f"Metadata: {dep.metadata}")- 带参数调用(传入
completed_tasks列表):走IDependency接口兼容路径,直接检查from_task_id是否在已完成列表中(task_star_line.py),适合星座级“就绪任务发现”场景; - 无参调用:返回内部
_is_satisfied标志,反映最近一次条件求值结果。
修改依赖
# Change dependency type dep.dependency_type = DependencyType.SUCCESS_ONLY # Update condition description dep.condition_description = "Updated: Deploy only after successful validation" # Set new condition evaluator def new_evaluator(result): return result.get("validation_score", 0) > 0.95 dep.set_condition_evaluator(new_evaluator) # Update metadata dep.update_metadata({ "updated_by": "admin", "reason": "Stricter validation threshold" })警告:执行期间修改需谨慎修改
dependency_type或condition_evaluator会重置满足状态——源码中两个 setter 均会清空_is_satisfied与_last_evaluation_result(task_star_line.py、L169-L180)。若在星座执行过程中修改依赖,可能导致已放行的任务重新被阻塞,务必谨慎。
metadata属性返回的是内部字典的副本(self._metadata.copy()),防止外部直接篡改内部状态;如需修改必须走update_metadata()并触发updated_at刷新。
序列化:JSON、字典与 Pydantic Schema
TaskStarLine 支持三种互为补充的序列化形态,便于持久化、传输与配置化加载。
JSON 导出 / 导入
# Export to JSON json_string = dep.to_json() print(json_string) # Save to file dep.to_json(save_path="dependency_backup.json") # Load from JSON string restored_dep = TaskStarLine.from_json(json_data=json_string) # Load from file loaded_dep = TaskStarLine.from_json(file_path="dependency_backup.json")to_json底层会先调用_ensure_json_serializable(task_star_line.py)清洗不可序列化字段:复杂对象转为vars()/字符串、集合转列表、可调用对象(如条件求值器)标记为<callable: 名称>占位符。因此导出的 JSON不含可执行的求值函数体——从 JSON 还原后,CONDITIONAL依赖会因缺少求值器而回退为成功判定语义,跨进程迁移时需注意这一点。
from_json要求json_data与file_path二选一且必须提供其一,二者同时或都不提供会抛出ValueError(task_star_line.py)。
字典转换
# Convert to dictionary dep_dict = dep.to_dict() # Create from dictionary new_dep = TaskStarLine.from_dict(dep_dict) # Dictionary structure print(dep_dict) # { # "line_id": "uuid-string", # "from_task_id": "task_a", # "to_task_id": "task_b", # "dependency_type": "success_only", # "condition_description": "...", # "metadata": {...}, # "is_satisfied": false, # "last_evaluation_result": null, # "created_at": "2025-11-06T...", # "updated_at": "2025-11-06T..." # }from_dict内部通过_parse_dependency_type(task_star_line.py)兼容字符串与枚举两种输入:字符串会做大小写不敏感映射,未知值回退为UNCONDITIONAL。同时它还会恢复is_satisfied、last_evaluation_result及三个时间戳字段,实现状态级还原。
Pydantic Schema 转换
# Convert to Pydantic BaseModel schema = dep.to_basemodel() # Create from Pydantic schema dep_from_schema = TaskStarLine.from_basemodel(schema)对应的 TaskStarLineSchema 定义了字段默认值(dependency_type默认"UNCONDITIONAL"、condition_description默认空串),并通过字段校验器将枚举值统一转为大写字符串(如DependencyType.SUCCESS_ONLY→"SUCCESS_ONLY"),通过模型校验器在缺省时自动生成line_id。from_basemodel会校验实例类型,非TaskStarLineSchema实例直接抛ValueError。
与 TaskConstellation 集成
添加依赖
from galaxy.constellation import TaskConstellation constellation = TaskConstellation(name="my_workflow") # Add tasks first constellation.add_task(task_a) constellation.add_task(task_b) # Add dependency try: constellation.add_dependency(dep) print("✅ Dependency added successfully") except ValueError as e: print(f"❌ Failed to add dependency: {e}")依赖校验
# TaskConstellation validates dependencies automatically try: # This would fail if it creates a cycle constellation.add_dependency(cyclic_dep) except ValueError as e: print(f"Validation error: {e}") # Output: "Adding dependency would create a cycle" # Check DAG validity is_valid, errors = constellation.validate_dag() if not is_valid: for error in errors: print(f"❌ {error}")TaskConstellation.add_dependency(task_constellation.py)执行三层校验后才真正挂载边:
- 任务存在性:
from_task_id与to_task_id必须都已注册,否则抛出ValueError(如"Source task xxx not found"); - 无环校验:调用
_would_create_cycle(task_constellation.py)——该函数以 DFS 检查to_task_id到from_task_id之间是否已存在可达路径,若存在则说明新增边会形成环,抛出"Adding dependency X -> Y would create a cycle"; - 引用同步:挂载边后同步更新两端 TaskStar 的依赖/被依赖引用(
from_task.add_dependent(...)/to_task.add_dependency(...)),并刷新星座状态。
validate_dag(task_constellation.py)则在整图层面检查,通过拓扑排序(get_topological_order)探测环,存在环时返回("DAG contains cycles", ...)错误列表。
高级模式
条件错误处理
将SUCCESS_ONLY与CONDITIONAL组合,即可实现"成功走 A 分支、失败走 B 分支"的经典 try/except 式工作流:
# Main task main_task = TaskStar( task_id="main_process", description="Process data" ) # Success path success_task = TaskStar( task_id="success_notification", description="Send success notification" ) # Error path error_task = TaskStar( task_id="error_recovery", description="Attempt recovery" ) # Success-only dependency success_dep = TaskStarLine.create_success_only( from_task_id="main_process", to_task_id="success_notification" ) # Failure-only dependency (using conditional) def on_failure(result): return result is None # Task failed if result is None failure_dep = TaskStarLine.create_conditional( from_task_id="main_process", to_task_id="error_recovery", condition_description="Run recovery if main task fails", condition_evaluator=on_failure )这里巧妙地利用了“成功 = 结果非None”的约定:on_failure恰好取反,让失败也能成为可路由的信号。
基于性能的路由
根据结果数据量在 GPU 与 CPU 两条处理路径间分流:
# Route to different processing paths based on data size def route_large_dataset(result): data_size = result.get("row_count", 0) return data_size > 1_000_000 # Route to GPU if > 1M rows # Route to GPU for large datasets gpu_dep = TaskStarLine.create_conditional( from_task_id="analyze_dataset", to_task_id="process_on_gpu", condition_description="Use GPU for datasets > 1M rows", condition_evaluator=route_large_dataset ) # Route to CPU for small datasets def route_small_dataset(result): data_size = result.get("row_count", 0) return data_size <= 1_000_000 cpu_dep = TaskStarLine.create_conditional( from_task_id="analyze_dataset", to_task_id="process_on_cpu", condition_description="Use CPU for datasets <= 1M rows", condition_evaluator=route_small_dataset )注意两个求值器互斥且互补,正好覆盖全量结果空间,形成稳定的分流决策。
错误处理
校验失败
# TaskStarLine validates on creation try: invalid_dep = TaskStarLine( from_task_id="task_a", to_task_id="task_a", # Self-loop! dependency_type=DependencyType.UNCONDITIONAL ) constellation.add_dependency(invalid_dep) except ValueError as e: print(f"Validation error: {e}") # TaskConstellation will detect cycle自环(self-loop)是_would_create_cycle的第一个捕获对象:DFS 检查to_task_id == from_task_id时立即返回True(task_constellation.py),因此在add_dependency阶段即被拦截。
求值器异常
def risky_evaluator(result): # This might raise an exception return result["complex_calculation"] / result["divisor"] dep = TaskStarLine.create_conditional( from_task_id="task_a", to_task_id="task_b", condition_description="Conditional with potential error", condition_evaluator=risky_evaluator ) # evaluate_condition catches exceptions and returns False result = {"complex_calculation": 100} # Missing "divisor" is_satisfied = dep.evaluate_condition(result) print(is_satisfied) # False (evaluator raised KeyError, caught internally) print(dep.last_evaluation_result) # False异常被evaluate_condition内部的except Exception捕获,返回False并保持_last_evaluation_result=False。这确保了单条依赖的求值失败只会阻塞依赖方,而不会让整个星座崩溃。
示例工作流
构建流水线(线性依赖)
# checkout → build → test → deploy checkout = TaskStar(task_id="checkout", description="Checkout code") build = TaskStar(task_id="build", description="Build project") test = TaskStar(task_id="test", description="Run tests") deploy = TaskStar(task_id="deploy", description="Deploy to production") # Sequential success-only dependencies dep1 = TaskStarLine.create_success_only("checkout", "build") dep2 = TaskStarLine.create_success_only("build", "test") dep3 = TaskStarLine.create_success_only("test", "deploy")任意一环失败(返回None),下游即被阻断,天然形成“质量门禁”。
扇出模式(Fan-Out)
# analyze → [process_gpu, process_cpu, process_edge] analyze = TaskStar(task_id="analyze", description="Analyze data") process_gpu = TaskStar(task_id="gpu", description="Process on GPU") process_cpu = TaskStar(task_id="cpu", description="Process on CPU") process_edge = TaskStar(task_id="edge", description="Process on edge device") # All three can start after analyze completes dep1 = TaskStarLine.create_unconditional("analyze", "gpu") dep2 = TaskStarLine.create_unconditional("analyze", "cpu") dep3 = TaskStarLine.create_unconditional("analyze", "edge")一个前置任务同时驱动多个并行下游,是 DAG 并行度的主要来源。
扇入模式(Fan-In)
# [task_a, task_b, task_c] → aggregate task_a = TaskStar(task_id="task_a", description="Process batch A") task_b = TaskStar(task_id="task_b", description="Process batch B") task_c = TaskStar(task_id="task_c", description="Process batch C") aggregate = TaskStar(task_id="aggregate", description="Aggregate results") # Aggregate waits for all three to complete dep1 = TaskStarLine.create_success_only("task_a", "aggregate") dep2 = TaskStarLine.create_success_only("task_b", "aggregate") dep3 = TaskStarLine.create_success_only("task_c", "aggregate")聚合任务会等待所有入边满足后才被调度,是 Map-Reduce 式工作流的天然载体。上述模式的正确性依赖TaskConstellation的get_ready_tasks调度逻辑(见 IDependencyResolver 与 task_constellation.py),该逻辑按“入边是否全部满足”决定任务是否可执行。
最佳实践
依赖设计准则
- 选对类型:根据工作流逻辑选择最贴切的依赖类型——顺序排队用
UNCONDITIONAL,质量门禁用SUCCESS_ONLY,收尾兜底用COMPLETION_ONLY,结果分流用CONDITIONAL; - 保持求值器简单:条件求值器应快速且确定性(deterministic),避免引入随机性或时间依赖;
- 处理好求值器异常:虽然
evaluate_condition内部会捕获异常并返回False,但异常会污染last_evaluation_result语义(失败与条件不满足难以区分),务必让求值器自身做防御性处理; - 写清条件描述:使用明确的
condition_description,它是调试星座执行轨迹时的第一手线索; - 避免环:
TaskConstellation会做无环校验,但设计阶段就应规划好层级,减少无效尝试。
好的求值器 vs 坏的求值器
✅好:简单、快速、防御式
def check_success(result): return result is not None and result.get("status") == "success"❌坏:复杂、缓慢、易错
def check_success(result): # Slow database query db_status = query_database(result["task_id"]) # Complex logic with potential errors return eval(result["complex_expression"]) and db_status常见陷阱
- 循环依赖:执行前务必校验 DAG 无环(
validate_dag);- 缺失任务:确保
from_task_id与to_task_id都已加入星座,否则add_dependency会抛ValueError;- 有状态求值器:避免依赖外部可变状态的求值器——同一结果两次求值可能得出不同结论;
- 慢求值器:求值在调度路径上同步执行,避免在求值器中做 I/O 或重计算。
相关组件
- TaskStar—— 原子执行单元,TaskStarLine 连接的对象;
- TaskConstellation—— DAG 管理器,负责校验并执行依赖;
- ConstellationEditor—— 支持撤销/重做的安全依赖编辑;
- Overview—— 任务星座框架总览。
API 参考
构造函数
TaskStarLine( from_task_id: str, to_task_id: str, dependency_type: DependencyType = DependencyType.UNCONDITIONAL, condition_description: Optional[str] = None, condition_evaluator: Optional[Callable[[Any], bool]] = None, line_id: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None )工厂方法
| 方法 | 说明 |
|---|---|
create_unconditional(from_id, to_id, desc) | 创建无条件依赖(类方法) |
create_success_only(from_id, to_id, desc) | 创建仅成功依赖(类方法) |
create_conditional(from_id, to_id, desc, evaluator) | 创建条件依赖(类方法) |
关键方法
| 方法 | 说明 |
|---|---|
evaluate_condition(result) | 求值条件是否满足(返回bool) |
mark_satisfied() | 手动标记为已满足 |
reset_satisfaction() | 重置满足状态 |
is_satisfied(completed_tasks=None) | 检查依赖是否满足;带参数走IDependency接口检查源任务是否完成,不带参数返回内部状态 |
set_condition_evaluator(evaluator) | 设置新的条件求值器 |
update_metadata(metadata) | 合并更新元数据 |
to_dict() | 转为字典 |
to_json(save_path) | 导出 JSON(可落盘) |
from_dict(data) | 从字典创建(类方法) |
from_json(json_data, file_path) | 从 JSON 字符串/文件创建(类方法) |
to_basemodel() | 转为 PydanticTaskStarLineSchema |
from_basemodel(schema) | 从 Pydantic schema 创建(类方法) |
代码定位
- 核心实现:galaxy/constellation/task_star_line.py
- 依赖类型枚举:galaxy/constellation/enums.py
- 接口契约
IDependency:galaxy/core/interfaces.py - Pydantic Schema:galaxy/agents/schema.py
- 星座级校验与调度:galaxy/constellation/task_constellation.py
- 相关测试:
tests/visualization/test_dependency_property_changes.py、tests/unit/schema/test_basemodel_integration.py、tests/examples/auto_id_example.py等覆盖了依赖属性变更、Schema 往返与自动 ID 分配等场景。
TaskStarLine —— 用智能依赖逻辑连接任务,让星座编排既有顺序的严谨,又有条件的灵活。
【免费下载链接】UFOUFO³: Weaving the Digital Agent Galaxy项目地址: https://gitcode.com/GitHub_Trending/uf/UFO
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考