-
-
Notifications
You must be signed in to change notification settings - Fork 0
14 task chain
Zuko edited this page Jan 24, 2026
·
2 revisions
Execute tasks sequentially với shared context và retry behaviors
TaskChain cho phép:
- Sequential task execution
- Shared context (ChainContext)
- Per-task retry behaviors
- Error handling strategies
- Error handling strategies
- Progress aggregation
-
Auto-tagging: Automaticaly tags children with
_ChainedChildandParent_{UUID}
from core.taskSystem import TaskChain, ChainRetryBehavior
chain = TaskChain(
name='Data Pipeline',
tasks=[task1, task2, task3],
description='Process data pipeline',
retryBehaviorMap={
'Task1': ChainRetryBehavior.RETRY_TASK,
'Task2': ChainRetryBehavior.SKIP_TASK,
'Task3': ChainRetryBehavior.FAIL_CHAIN
}
)taskManager.addChainTask(
name='Pipeline',
tasks=[task1, task2, task3],
retryBehaviorMap={...}
)Shared state between tasks:
# In task
def handle(self):
# Access chain context
context = self._chainContext
# Set data
context.set('user_id', 123)
context.set('data', {'key': 'value'})
# Get data
userId = context.get('user_id')
data = context.get('data', default={})from core.taskSystem import ChainRetryBehavior
ChainRetryBehavior.RETRY_TASK # Retry failed task
ChainRetryBehavior.SKIP_TASK # Skip and continue
ChainRetryBehavior.FAIL_CHAIN # Fail entire chainfrom core.taskSystem import AbstractTask
class FetchTask(AbstractTask):
def handle(self):
import requests
response = requests.get('https://api.example.com/data')
# Share data with next task
self._chainContext.set('raw_data', response.json())
class ProcessTask(AbstractTask):
def handle(self):
# Get data from previous task
rawData = self._chainContext.get('raw_data')
# Process...
processed = self._process(rawData)
# Share with next task
self._chainContext.set('processed_data', processed)
class SaveTask(AbstractTask):
def handle(self):
# Get processed data
data = self._chainContext.get('processed_data')
# Save to database
self._save(data)
# Create chain
chain = taskManager.addChainTask(
name='Data Pipeline',
tasks=[FetchTask(), ProcessTask(), SaveTask()]
)chain = taskManager.addChainTask(
name='Robust Pipeline',
tasks=[
FetchTask(name='Fetch'),
ProcessTask(name='Process'),
SaveTask(name='Save')
],
retryBehaviorMap={
'Fetch': ChainRetryBehavior.RETRY_TASK, # Retry on network error
'Process': ChainRetryBehavior.SKIP_TASK, # Skip if processing fails
'Save': ChainRetryBehavior.FAIL_CHAIN # Critical - fail chain
}
)class LoginTask(AbstractTask):
def handle(self):
ctx = QtAppContext.globalInstance()
browser = ctx.getService('browser')
browser.navigate('https://example.com/login')
browser.fillForm({'username': 'user', 'password': 'pass'})
browser.click('#login-button')
# Share session
self._chainContext.set('logged_in', True)
class ScrapeTask(AbstractTask):
def handle(self):
if not self._chainContext.get('logged_in'):
self.fail('Not logged in')
return
ctx = QtAppContext.globalInstance()
browser = ctx.getService('browser')
data = browser.extractData('.data-table')
self._chainContext.set('scraped_data', data)
class ExportTask(AbstractTask):
def handle(self):
data = self._chainContext.get('scraped_data')
with open('output.json', 'w') as f:
json.dump(data, f)
# Chain with shared browser
chain = taskManager.addChainTask(
name='Scraping Pipeline',
tasks=[LoginTask(), ScrapeTask(), ExportTask()],
retryBehaviorMap={
'LoginTask': ChainRetryBehavior.RETRY_TASK,
'ScrapeTask': ChainRetryBehavior.RETRY_TASK,
'ExportTask': ChainRetryBehavior.FAIL_CHAIN
}
)# Chain aggregates progress from all tasks
chain.progressUpdated.connect(lambda progress: print(f'Chain: {progress}%'))
# Individual task progress
task1.progressUpdated.connect(lambda progress: print(f'Task 1: {progress}%'))# Task fails
class RiskyTask(AbstractTask):
def handle(self):
if error_condition:
self.fail('Something went wrong')
# Chain behavior depends on retryBehaviorMap
chain = taskManager.addChainTask(
name='Chain',
tasks=[RiskyTask()],
retryBehaviorMap={
'RiskyTask': ChainRetryBehavior.SKIP_TASK # Continue despite failure
}
)# Use descriptive task names
tasks = [
FetchDataTask(name='FetchData'),
ProcessDataTask(name='ProcessData'),
SaveDataTask(name='SaveData')
]
# Share data via context
self._chainContext.set('key', value)
# Check context data exists
data = self._chainContext.get('key')
if data is None:
self.fail('Missing required data')
return
# Use appropriate retry behaviors
retryBehaviorMap={
'NetworkTask': ChainRetryBehavior.RETRY_TASK,
'OptionalTask': ChainRetryBehavior.SKIP_TASK,
'CriticalTask': ChainRetryBehavior.FAIL_CHAIN
}# Don't use global variables
global_data = None # Wrong! Use context
def handle(self):
global global_data
global_data = result # Use self._chainContext.set()
# Don't assume task order
def handle(self):
data = self._chainContext.get('data')
# What if previous task was skipped?
# Don't forget error handling
def handle(self):
data = self._chainContext.get('required_data')
process(data) # May be None!- AbstractTask - Base task class
- TaskManager - Chain creation
- Task System Overview - Architecture