111 lines
3.6 KiB
Python
111 lines
3.6 KiB
Python
"""Domain-owned AI capabilities; caller owns authorization and transactions."""
|
||
|
||
from fastapi import HTTPException
|
||
from pydantic import Field
|
||
|
||
from ..schemas import Contract, JobInput
|
||
from .capabilities import Capability, EmptyArgs
|
||
|
||
|
||
class JobArgs(Contract):
|
||
job_id: str = Field(min_length=1, max_length=100)
|
||
|
||
|
||
async def list_jobs(ctx, args):
|
||
return {"items": (await ctx.business.list_jobs())[:20]}
|
||
|
||
|
||
async def preview_create(ctx, args):
|
||
return {"operation": args.model_dump(mode="json")}
|
||
|
||
|
||
async def preview_job(ctx, args):
|
||
return {"job": await ctx.business.get_job_status(args.job_id)}
|
||
|
||
|
||
async def check_job(ctx, args, preview):
|
||
current = await ctx.business.get_job_status(args.job_id)
|
||
if current["updated_at"] != preview["job"]["updated_at"] or current["status"] != preview["job"]["status"]:
|
||
raise HTTPException(409, "任务状态已变化,请重新确认操作")
|
||
|
||
|
||
async def cancel_job(ctx, args, preview):
|
||
await check_job(ctx, args, preview)
|
||
return await ctx.business.cancel_job(args.job_id)
|
||
|
||
|
||
async def retry_job(ctx, args, preview):
|
||
await check_job(ctx, args, preview)
|
||
return await ctx.business.retry_job(args.job_id)
|
||
|
||
|
||
async def wake_sync(runner, result):
|
||
"""Called only after the job and audit commit."""
|
||
runner.wake.set()
|
||
|
||
|
||
async def cancel_sync(runner, result):
|
||
"""The durable cancel decision precedes interruption of the in-process task."""
|
||
await runner.cancel(result["job_id"])
|
||
|
||
|
||
INSTRUCTIONS = "任务创建后返回任务信息并结束本轮,不要循环等待任务完成。不推测未执行操作已经成功。"
|
||
|
||
|
||
CAPABILITIES = (
|
||
Capability(
|
||
name="list_jobs",
|
||
schema=EmptyArgs,
|
||
description="查询最近的同步任务,不要循环轮询等待。",
|
||
label="查询任务",
|
||
renderer="jobs",
|
||
effect="query",
|
||
handler=list_jobs,
|
||
),
|
||
Capability(
|
||
name="get_job_status",
|
||
schema=JobArgs,
|
||
description="查询指定任务的状态、目标和错误,不要循环等待任务完成。",
|
||
label="查看任务状态",
|
||
renderer="jobs",
|
||
effect="query",
|
||
handler=lambda ctx, args: ctx.business.get_job_status(args.job_id),
|
||
),
|
||
Capability(
|
||
name="create_sync_job",
|
||
schema=JobInput,
|
||
description="提出同步或本地自相关任务,等待确认。full_sync 仅同步已提交;待提交必须用 daily_sync 并指定 submission、date_from/date_to(UTC),待提交按创建日、已提交按提交日逐天同步。alpha_refresh/pnl_refresh/self_correlation 使用固定 alpha_ids;自相关缺失 PnL 时自动补取,不触发平台检查。创建后立即返回任务 ID。",
|
||
label="创建同步任务",
|
||
renderer="jobs",
|
||
effect="confirm",
|
||
preview=preview_create,
|
||
execute=lambda ctx, args, preview: ctx.business.create_sync_job(args),
|
||
after_commit=wake_sync,
|
||
refresh=("jobs",),
|
||
),
|
||
Capability(
|
||
name="cancel_job",
|
||
schema=JobArgs,
|
||
description="提出取消指定同步任务,等待用户确认。",
|
||
label="取消任务",
|
||
renderer="jobs",
|
||
effect="confirm",
|
||
preview=preview_job,
|
||
execute=cancel_job,
|
||
after_commit=cancel_sync,
|
||
refresh=("jobs",),
|
||
),
|
||
Capability(
|
||
name="retry_job",
|
||
schema=JobArgs,
|
||
description="提出重试指定失败或暂停的同步任务,等待用户确认。",
|
||
label="重试任务",
|
||
renderer="jobs",
|
||
effect="confirm",
|
||
preview=preview_job,
|
||
execute=retry_job,
|
||
after_commit=wake_sync,
|
||
refresh=("jobs",),
|
||
),
|
||
)
|