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",),
|
|||
|
|
),
|
|||
|
|
)
|