Compare commits

..

9 Commits

Author SHA1 Message Date
k88936 3467042000 Merge branch 'master' into feat/fail-as-event
Project CI / AI game creator shell Rust smoke (pull_request) Has been cancelled
Project CI / AI game creator shell Rust crates (pull_request) Has been cancelled
Project CI / Backend tests (pull_request) Has been cancelled
Project CI / Native shell tests (pull_request) Has been cancelled
Project CI / Frontend tests (pull_request) Has been cancelled
Project CI / Repository checks (pull_request) Has been cancelled
Project CI / AI game creator shell web tests (pull_request) Has been cancelled
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Has been cancelled
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Has been cancelled
2026-09-22 20:19:23 +08:00
k88936 80b15b24ae appSurface 用例:失败只经终态事件收口,界面不再停在"还在处理"
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Has been cancelled
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Has been cancelled
Project CI / AI game creator shell Rust smoke (pull_request) Has been cancelled
Project CI / AI game creator shell Rust crates (pull_request) Has been cancelled
Project CI / Backend tests (pull_request) Has been cancelled
Project CI / Native shell tests (pull_request) Has been cancelled
Project CI / Frontend tests (pull_request) Has been cancelled
Project CI / Repository checks (pull_request) Has been cancelled
Project CI / AI game creator shell web tests (pull_request) Has been cancelled
- 新增「closes the turn from the host failure payload instead of leaving it running」:宿主发过 turn.started 之后以 failure 载荷收场并让命令失败,断言失败文案来自事件载荷(经同一份可见文案映射)、"陶泥儿正在处理"消失、终止钮消失、输入盒回到「发送」
- 同一条用例反向断言命令返回的错误原文不进聊天:那条通道只负责运行错误横幅
- 变异校验:让 reducer 不落失败说明条目时该用例变红(1 failed),恢复后绿
2026-09-22 17:40:17 +08:00
k88936 b432556ee3 reducer 用例:失败终态落说明条目、文案与横幅同源、重放不重复
- 新增「失败终态(turn.completed 带 failure 载荷)」六条用例:失败照样收口并冻结终点、说明条目按本轮开口身份派生、本轮开口条目拿到边界
- 可见文案与运行错误横幅共用同一份映射(点名 codex-app-server-error:context-window-exceeded)
- 重复 / 迟到的失败终态不追加第二条说明、不抬高冻结终点、不复活运行态
- 身份不匹配的失败终态不动正在跑的这一轮;空原因不落说明条目但终态照样收口
- 没有身份时用事件时间派生说明身份,两轮失败不会合并成一条
2026-09-22 17:40:17 +08:00
k88936 258645f182 前端失败说明改由事件驱动:turn.completed.failure 落成本轮说明条目,命令返回只留横幅
- 新增 conversation/directTurnFailure.ts:失败说明条目的展示身份(本轮开口身份 + :failure)与可见文案(复用 projectRuntimeVisibleError)两条口径集中一处
- directThreadChat 的 turn.completed 分支读 failure 载荷:非空原因先落成本轮最后一条说明条目,再走同一个收口函数;失败不再是第二套生命周期
- useDirectProjectChatController 的失败分支不再写聊天气泡:聊天文案唯一来源是事件,命令返回只保留运行错误横幅(含详情 long detail)与诊断留痕
- 数据流、时序与投影注释同步:标注失败说明来自事件、本地通道只剩终止说明与壳层 announce
2026-09-22 17:40:17 +08:00
k88936 5cf4a018b3 宿主终态接线:失败写进 turn.completed 的 failure 载荷,并在 turn.started 之后武装兜底守卫
- codex_app_server 的 DirectProject 终态改用 direct_turn_failure 判定:失败走 turn_completed_failed(原因脱敏 + 截断后写进同一个事件),其余仍走 turn_completed(status)
- 失败判定两个来源:collect_result 是 Err 时用错误本身当原因;collect_result 是交付报告但 status 已判成 failed 时用那份报告当原因
- turn.started 进入队列后立即武装 DirectTurnFailureGuard,写完终态 disarm:panic、回合 future 被丢弃、终态之前的早退都会补一条 host-dropped 失败终态,前端不会停在"还在跑"
- 定向 `cargo test direct_`(438 passed,含 wire / manager / 失败策略模块)
2026-09-22 17:40:17 +08:00
k88936 5d9223c32e 失败终态策略独立成模块:分类、原因脱敏与 Drop 兜底守卫
- 新增 agent/direct_turn_failure.rs:LlmError → 稳定分类(timeout / model-failed / transport-failed / request-rejected)、判定"终态是不是失败"并给出脱敏截断后的原因、DirectTurnFailureGuard(turn.started 之后武装、写完终态 disarm,Drop 时补 host-dropped 失败终态)
- 守卫兜底覆盖 panic / future 被丢弃 / 终态之前的早退;kill -9 与 turn.started 之前的早退写进模块注释,明确不为它们补路径
- agent.rs 注册模块并再导出
- 5 条用例:错误分类映射、只有 failed 终态带载荷、原因脱敏 + 按字符截断、armed 后 Drop 补终态、disarm 后不再产出事件
2026-09-22 17:40:17 +08:00
k88936 d32c99c927 事件协议:turn.completed 增加可选 failure 载荷,失败终态与正常终态同权入锚点
- direct_thread_wire 新增 DirectTurnFailure{kind,message} 类型,给 TurnCompleted 增可选 failure 字段,并补 turn_completed_failed 构造器与 failure 读取器
- with_user_item_id 显式带上 failure:原先把 TurnCompleted 写成 `..` 会静默吞掉失败载荷,身份与原因必须一起流转
- 新增 wire 用例:失败终态带载荷、正常终态不带且回写不补 null、缺载荷的 failed 事件仍可反序列化
- direct_thread_manager 增回归用例:turn.completed(status=failed) 必须顶替更早的 turn.started 成为 lifecycle_anchor,重放不会把已收口的回合看成"还在跑"
- 重新生成 ts-rs 绑定(新增 DirectTurnFailure.ts、DirectThreadEvent.ts 增 failure 字段)并按 prettier 格式化
2026-09-22 17:40:17 +08:00
k88936 d36f5842b6 文档:失败回合终态定为 turn.completed 带 failure 载荷,宿主 Drop 守卫兜底
- ADR【DirectProject对话历史单一事实源】补三条决策:终态事件只有 turn.completed,失败时 status="failed" 必须带 failure{kind,message};宿主 Drop 守卫在 turn.started 之后武装、写完终态即解除;失败原因只走事件这条通道,聊天说明的展示位保留、数据来源换成事件
- 同 ADR「影响」补两条已知边界(进程被强杀时没有 Drop、turn.started 之前的早退不产回合也不补终态)与「可见文案映射规则不变」的口径
- 技术方案【DirectProject Codex原始历史与异常恢复】同步线上形状:turn.completed 增加可选 failure,并写明失败终态与正常终态同权顶替 lifecycle_anchor
- decision-log 记本次决策、明确不做项、影响范围与验证方式
2026-09-22 17:40:17 +08:00
k88936 981a6b0021 Revert:撤掉前端「命令返回就收口」的兜底,改由宿主 turn.failed 事件收口
- 撤销 a35956f3e 的前端实现:directThreadChat 的 commandClosedTurnUserItemId / stopDirectThreadTurn、subscription 的 stopCommandTurn、controller 失败分支的调用,以及随附的两处用例
- 原因:失败语义改由宿主事件(turn.failed)表达,前端不再自造第二条「结束」判定路径,也不再在失败路径上补本地消息
2026-09-22 17:40:16 +08:00
521 changed files with 29850 additions and 16604 deletions
-6
View File
@@ -239,12 +239,6 @@ GENARRATIVE_ENABLE_IMAGE_EDITOR_AGENT_SIDEBAR="false"
# Windows/macOS 是系统维度,不填写 dev-win/dev-mac。 # Windows/macOS 是系统维度,不填写 dev-win/dev-mac。
GENARRATIVE_CLIENT_DOWNLOAD_CHANNEL="dev" GENARRATIVE_CLIENT_DOWNLOAD_CHANNEL="dev"
# 客户端埋点接收绑定的公开 origin,由 API Server 运行时读取;修改后重启服务。
# 必须与客户端登录地址一致,不带 /api、路径或尾部斜杠;未配置/非法时上传接口返回 503。
# 本地端口按实际启动结果填写(端口漂移后需同步),localhost 与 127.0.0.1 不可混用。
# dev 使用 https://dev.genarrative.worldrelease 使用 https://www.genarrative.world。
GENARRATIVE_AGC_ANALYTICS_ORIGIN="http://127.0.0.1:8082"
# Optional: official VikingDB credentials for regenerating build-tag similarities # Optional: official VikingDB credentials for regenerating build-tag similarities
# with the Python embedding script. The script auto-loads `.env.local` and uses # with the Python embedding script. The script auto-loads `.env.local` and uses
# the fixed `bge-large-zh` embedding model. # the fixed `bge-large-zh` embedding model.
@@ -7,7 +7,6 @@ import {
getAdminFeatureGateConfig, getAdminFeatureGateConfig,
getAdminUserDetail, getAdminUserDetail,
importAdminAgcTemplates, importAdminAgcTemplates,
listAdminAgcTrackingEvents,
listAdminGameDistributionReviews, listAdminGameDistributionReviews,
listAdminRechargeOrders, listAdminRechargeOrders,
reconcileAdminUserConsumption, reconcileAdminUserConsumption,
@@ -25,30 +24,6 @@ afterEach(() => {
vi.unstubAllGlobals(); vi.unstubAllGlobals();
}); });
test('客户端埋点查询传递筛选和游标并复用后台认证', async () => {
const payload = { entries: [], nextCursor: null };
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ ok: true, data: payload }), {
status: 200,
}),
);
vi.stubGlobal('fetch', fetchMock);
expect(
await listAdminAgcTrackingEvents('token', {
userId: 'user+1',
projectId: 'project-1',
cursor: 'page/2',
limit: 50,
}),
).toEqual(payload);
expect(fetchMock).toHaveBeenCalledWith(
'/admin/api/agc/tracking-events?userId=user%2B1&projectId=project-1&cursor=page%2F2&limit=50',
expect.objectContaining({
headers: expect.objectContaining({ Authorization: 'Bearer token' }),
}),
);
});
test('模板管理读取和更新复用认证封装,提交 revision 和封面但不提交 ZIP 或版本', async () => { test('模板管理读取和更新复用认证封装,提交 revision 和封面但不提交 ZIP 或版本', async () => {
const library = { revision: 'revision-new', writable: true, templates: [] }; const library = { revision: 'revision-new', writable: true, templates: [] };
const fetchMock = vi.fn().mockImplementation( const fetchMock = vi.fn().mockImplementation(
-27
View File
@@ -1,8 +1,6 @@
import type { import type {
AdminAccountListResponse, AdminAccountListResponse,
AdminAgcTemplateLibraryResponse, AdminAgcTemplateLibraryResponse,
AdminAgcTrackingEventListResponse,
AdminAgcTrackingEventQuery,
AdminConfirmEditorShowcaseCampaignImageUploadRequest, AdminConfirmEditorShowcaseCampaignImageUploadRequest,
AdminCreateAccountRequest, AdminCreateAccountRequest,
AdminCreateAccountResponse, AdminCreateAccountResponse,
@@ -411,31 +409,6 @@ export function listAdminTrackingEventKeys(token: string) {
); );
} }
export function listAdminAgcTrackingEvents(
token: string,
query: AdminAgcTrackingEventQuery = {},
) {
return request<AdminAgcTrackingEventListResponse>(
`/admin/api/agc/tracking-events${buildQueryString((params) => {
for (const key of [
'userId',
'projectId',
'creativeTaskId',
'agentRunId',
'eventName',
'clientVersion',
'startTime',
'endTime',
'cursor',
] as const) {
appendQueryParam(params, key, query[key]);
}
appendNumericQueryParam(params, 'limit', query.limit);
})}`,
{ token },
);
}
export function listAdminErrorReports( export function listAdminErrorReports(
token: string, token: string,
query: { query: {
-38
View File
@@ -820,44 +820,6 @@ export interface AdminTrackingEventListResponse {
entries: AdminTrackingEventEntryPayload[]; entries: AdminTrackingEventEntryPayload[];
} }
export interface AdminAgcTrackingEventQuery {
userId?: string;
projectId?: string;
creativeTaskId?: string;
agentRunId?: string;
eventName?: string;
clientVersion?: string;
startTime?: string;
endTime?: string;
cursor?: string;
limit?: number;
}
export interface AdminAgcTrackingEventEntry {
eventId: string;
schemaVersion: number;
eventName: string;
eventTime: string;
userId: string;
editorSessionId: string;
projectId: string | null;
creativeTaskId: string | null;
agentRunId: string | null;
agentTurnId: string | null;
status: string | null;
errorCode: string | null;
source: string;
clientVersion: string;
properties: Record<string, unknown>;
batchId: string;
receivedAt: string;
}
export interface AdminAgcTrackingEventListResponse {
entries: AdminAgcTrackingEventEntry[];
nextCursor: string | null;
}
export interface AdminTrackingEventKeyPayload { export interface AdminTrackingEventKeyPayload {
eventKey: string; eventKey: string;
eventTitle: string; eventTitle: string;
-7
View File
@@ -20,7 +20,6 @@ import {
import { AdminAccountsPage } from '../pages/AdminAccountsPage'; import { AdminAccountsPage } from '../pages/AdminAccountsPage';
import { AdminAgcModelsPage } from '../pages/AdminAgcModelsPage'; import { AdminAgcModelsPage } from '../pages/AdminAgcModelsPage';
import { AdminAgcTemplatesPage } from '../pages/AdminAgcTemplatesPage'; import { AdminAgcTemplatesPage } from '../pages/AdminAgcTemplatesPage';
import { AdminAgcTrackingPage } from '../pages/AdminAgcTrackingPage';
import { AdminDashboardPage } from '../pages/AdminDashboardPage'; import { AdminDashboardPage } from '../pages/AdminDashboardPage';
import { AdminDatabaseTablesPage } from '../pages/AdminDatabaseTablesPage'; import { AdminDatabaseTablesPage } from '../pages/AdminDatabaseTablesPage';
import { AdminDebugHttpPage } from '../pages/AdminDebugHttpPage'; import { AdminDebugHttpPage } from '../pages/AdminDebugHttpPage';
@@ -234,12 +233,6 @@ export function AdminApp() {
onUnauthorized={handleUnauthorized} onUnauthorized={handleUnauthorized}
/> />
) : null} ) : null}
{activeRouteId === 'agc-tracking' ? (
<AdminAgcTrackingPage
token={token}
onUnauthorized={handleUnauthorized}
/>
) : null}
{activeRouteId === 'error-reports' ? ( {activeRouteId === 'error-reports' ? (
<AdminErrorReportsPage <AdminErrorReportsPage
token={token} token={token}
-1
View File
@@ -40,7 +40,6 @@ const routeIcons = {
tables: Database, tables: Database,
debug: Bug, debug: Bug,
tracking: Table2, tracking: Table2,
'agc-tracking': Table2,
'error-reports': Bug, 'error-reports': Bug,
'gray-release': GitBranch, 'gray-release': GitBranch,
redeem: TicketPercent, redeem: TicketPercent,
@@ -8,22 +8,6 @@ import {
routeHash, routeHash,
} from './adminRoutes'; } from './adminRoutes';
test('客户端埋点路由遵守成员页签权限', () => {
expect(resolveAdminRoute('#agc-tracking')).toBe('agc-tracking');
expect(
getAccessibleAdminRoutes({
accountRole: 'member',
tabPermissions: ['agc-tracking'],
}).map((route) => route.id),
).toEqual(['agc-tracking']);
expect(
getAccessibleAdminRoutes({
accountRole: 'member',
tabPermissions: ['tracking'],
}).some((route) => route.id === 'agc-tracking'),
).toBe(false);
});
test('后台默认进入 Dashboard', () => { test('后台默认进入 Dashboard', () => {
expect(adminRoutes[0]).toEqual({ expect(adminRoutes[0]).toEqual({
id: 'dashboard', id: 'dashboard',
-2
View File
@@ -5,7 +5,6 @@ export type AdminRouteId =
| 'tables' | 'tables'
| 'debug' | 'debug'
| 'tracking' | 'tracking'
| 'agc-tracking'
| 'error-reports' | 'error-reports'
| 'gray-release' | 'gray-release'
| 'redeem' | 'redeem'
@@ -42,7 +41,6 @@ export const adminRoutes: AdminRouteDefinition[] = [
{ id: 'tables', label: '表查询', hash: '#tables' }, { id: 'tables', label: '表查询', hash: '#tables' },
{ id: 'debug', label: 'API 调试', hash: '#debug' }, { id: 'debug', label: 'API 调试', hash: '#debug' },
{ id: 'tracking', label: '埋点数据', hash: '#tracking' }, { id: 'tracking', label: '埋点数据', hash: '#tracking' },
{ id: 'agc-tracking', label: '客户端埋点', hash: '#agc-tracking' },
{ id: 'error-reports', label: '错误报告', hash: '#error-reports' }, { id: 'error-reports', label: '错误报告', hash: '#error-reports' },
{ id: 'gray-release', label: '灰度发布', hash: '#gray-release' }, { id: 'gray-release', label: '灰度发布', hash: '#gray-release' },
{ id: 'redeem', label: '兑换码', hash: '#redeem' }, { id: 'redeem', label: '兑换码', hash: '#redeem' },
@@ -1,129 +0,0 @@
// @vitest-environment jsdom
import {
cleanup,
fireEvent,
render,
screen,
waitFor,
within,
} from '@testing-library/react';
import { afterEach, beforeEach, expect, test, vi } from 'vitest';
import {
AdminApiError,
listAdminAgcTrackingEvents,
} from '../api/adminApiClient';
import type { AdminAgcTrackingEventEntry } from '../api/adminApiTypes';
import { AdminAgcTrackingPage } from './AdminAgcTrackingPage';
vi.mock('../api/adminApiClient', async (original) => ({
...(await original<typeof import('../api/adminApiClient')>()),
listAdminAgcTrackingEvents: vi.fn(),
}));
vi.mock('../components/AdminUserReferenceButton', () => ({
AdminUserReferenceButton: () => null,
}));
const entry: AdminAgcTrackingEventEntry = {
eventId: 'event-1',
schemaVersion: 1,
eventName: 'project_save',
eventTime: '2026-09-21T12:00:00.123Z',
userId: 'user-1',
editorSessionId: 'session-1',
projectId: 'project-1',
creativeTaskId: 'project-1',
agentRunId: null,
agentTurnId: null,
status: 'success',
errorCode: null,
source: 'gui',
clientVersion: '1.0',
properties: { reason: 'manual' },
batchId: 'batch-1',
receivedAt: '2026-09-21T12:15:00.000Z',
};
beforeEach(() => {
vi.mocked(listAdminAgcTrackingEvents).mockReset();
});
afterEach(cleanup);
test('翻页沿用游标,筛选和刷新从第一页重新查询,详情保留 null 和事件属性', async () => {
vi.mocked(listAdminAgcTrackingEvents)
.mockResolvedValueOnce({ entries: [entry], nextCursor: 'cursor-page-2' })
.mockResolvedValueOnce({
entries: [{ ...entry, eventId: 'event-2' }],
nextCursor: null,
})
.mockResolvedValue({ entries: [entry], nextCursor: 'cursor-new' });
render(<AdminAgcTrackingPage token="admin-token" onUnauthorized={vi.fn()} />);
await screen.findByText('保存项目', { selector: 'td' });
expect(listAdminAgcTrackingEvents).toHaveBeenLastCalledWith('admin-token', {
cursor: undefined,
limit: 50,
});
fireEvent.click(screen.getByText('查看详情'));
const dialog = screen.getByRole('dialog');
expect(within(dialog).getAllByText('—').length).toBe(3);
expect(within(dialog).getByText(/"reason": "manual"/)).toBeTruthy();
expect(within(dialog).getByText(entry.eventTime)).toBeTruthy();
fireEvent.click(within(dialog).getByLabelText('关闭详情'));
fireEvent.click(screen.getByText('下一页'));
await waitFor(() =>
expect(listAdminAgcTrackingEvents).toHaveBeenLastCalledWith('admin-token', {
cursor: 'cursor-page-2',
limit: 50,
}),
);
await waitFor(() =>
expect(screen.getByText('查询').closest('button')?.disabled).toBe(false),
);
fireEvent.click(screen.getByText('上一页'));
await screen.findByText('第 1 页');
expect(listAdminAgcTrackingEvents).toHaveBeenCalledTimes(2);
fireEvent.change(screen.getByLabelText('项目 ID'), {
target: { value: 'project-2' },
});
fireEvent.click(screen.getByText('查询'));
await waitFor(() =>
expect(listAdminAgcTrackingEvents).toHaveBeenLastCalledWith(
'admin-token',
expect.objectContaining({
projectId: 'project-2',
cursor: undefined,
limit: 50,
}),
),
);
expect(screen.getByText('第 1 页')).toBeTruthy();
await waitFor(() =>
expect(screen.getByText('刷新').closest('button')?.disabled).toBe(false),
);
fireEvent.click(screen.getByText('刷新'));
await waitFor(() =>
expect(listAdminAgcTrackingEvents).toHaveBeenCalledTimes(4),
);
});
test('空结果、权限失败与登录失效沿用后台反馈', async () => {
const unauthorized = vi.fn();
vi.mocked(listAdminAgcTrackingEvents)
.mockResolvedValueOnce({ entries: [], nextCursor: null })
.mockRejectedValueOnce(
new AdminApiError({ status: 403, message: '无权访问客户端埋点' }),
)
.mockRejectedValueOnce(
new AdminApiError({ status: 401, message: '登录失效' }),
);
render(<AdminAgcTrackingPage token="token" onUnauthorized={unauthorized} />);
await screen.findByText('暂无客户端埋点数据');
fireEvent.click(screen.getByText('刷新'));
expect((await screen.findByRole('alert')).textContent).toContain(
'无权访问客户端埋点',
);
fireEvent.click(screen.getByText('刷新'));
await waitFor(() =>
expect(unauthorized).toHaveBeenCalledWith('登录状态已失效'),
);
});
@@ -1,379 +0,0 @@
import { Modal } from '@genarrative/shared/components';
import { FormEvent, useEffect, useRef, useState } from 'react';
import { listAdminAgcTrackingEvents } from '../api/adminApiClient';
import type {
AdminAgcTrackingEventEntry,
AdminAgcTrackingEventListResponse,
AdminAgcTrackingEventQuery,
} from '../api/adminApiTypes';
import { AdminUserReferenceButton } from '../components/AdminUserReferenceButton';
import { handlePageError } from './pageUtils';
const eventLabels: Record<string, string> = {
editor_session_start: '编辑器会话开始',
editor_session_end: '编辑器会话结束',
editor_focus_start: '编辑器获得焦点',
editor_focus_end: '编辑器失去焦点',
project_create_success: '项目创建成功',
project_open: '打开项目',
creative_task_submit: '首次提交创作目标',
agent_run_completed: 'Agent 运行完成',
agent_run_failed: 'Agent 运行失败',
project_revision_created: '项目产生修改',
preview_ready: '预览就绪',
project_save: '保存项目',
};
const fieldLabels: Record<keyof AdminAgcTrackingEventEntry, string> = {
eventId: '事件 ID',
schemaVersion: '事件版本',
eventName: '事件类型',
eventTime: '发生时间',
userId: '用户 ID',
editorSessionId: '编辑器会话 ID',
projectId: '项目 ID',
creativeTaskId: '创作目标 ID',
agentRunId: 'Agent run ID',
agentTurnId: 'Agent turn ID',
status: '结果',
errorCode: '错误码',
source: '来源',
clientVersion: '客户端版本',
properties: '事件属性',
batchId: '批次 ID',
receivedAt: '入库时间',
};
function formatTime(value: string) {
const date = new Date(value);
return Number.isNaN(date.getTime()) ? value : date.toLocaleString('zh-CN');
}
export function AdminAgcTrackingPage({
token,
onUnauthorized,
}: {
token: string;
onUnauthorized: (message?: string) => void;
}) {
const [filters, setFilters] = useState({
userId: '',
projectId: '',
eventName: '',
startTime: '',
endTime: '',
});
const [query, setQuery] = useState<AdminAgcTrackingEventQuery>({});
const [cursors, setCursors] = useState<Array<string | undefined>>([
undefined,
]);
const [page, setPage] = useState(0);
const [refresh, setRefresh] = useState(0);
const [entries, setEntries] = useState<AdminAgcTrackingEventEntry[]>([]);
const [nextCursor, setNextCursor] = useState<string | null>(null);
const [loading, setLoading] = useState(true);
const [error, setError] = useState('');
const [detail, setDetail] = useState<AdminAgcTrackingEventEntry | null>(null);
const cursor = cursors[page];
// 首页没有入站游标,返回首页时保留原快照,刷新才获取新数据。
const firstPage = useRef<{
token: string;
query: AdminAgcTrackingEventQuery;
response: AdminAgcTrackingEventListResponse;
} | null>(null);
useEffect(() => {
const saved = firstPage.current;
if (!cursor && saved?.token === token && saved.query === query) {
setEntries(saved.response.entries);
setNextCursor(saved.response.nextCursor);
setError('');
setLoading(false);
return;
}
let active = true;
setLoading(true);
setError('');
setEntries([]);
setNextCursor(null);
void listAdminAgcTrackingEvents(token, { ...query, cursor, limit: 50 })
.then((response) => {
if (!active) return;
if (!cursor) firstPage.current = { token, query, response };
setEntries(response.entries);
setNextCursor(response.nextCursor);
})
.catch((failure: unknown) => {
if (active) handlePageError(failure, onUnauthorized, setError);
})
.finally(() => {
if (active) setLoading(false);
});
return () => {
active = false;
};
}, [token, query, cursor, refresh, onUnauthorized]);
function resetPages() {
firstPage.current = null;
setCursors([undefined]);
setPage(0);
setDetail(null);
}
function search(event: FormEvent<HTMLFormElement>) {
event.preventDefault();
if (
filters.startTime &&
filters.endTime &&
new Date(filters.startTime) >= new Date(filters.endTime)
) {
setError('结束时间必须晚于开始时间');
return;
}
resetPages();
setQuery({
userId: filters.userId.trim(),
projectId: filters.projectId.trim(),
eventName: filters.eventName,
startTime: filters.startTime
? new Date(filters.startTime).toISOString()
: undefined,
endTime: filters.endTime
? new Date(filters.endTime).toISOString()
: undefined,
});
}
function related(
key: 'creativeTaskId' | 'agentRunId' | 'clientVersion',
value: string,
) {
resetPages();
setFilters({
userId: '',
projectId: '',
eventName: '',
startTime: '',
endTime: '',
});
setQuery({ [key]: value });
}
return (
<section className="admin-page admin-page-wide">
<div className="admin-page-heading">
<div>
<h2></h2>
<p></p>
</div>
<button
type="button"
className="admin-secondary-button"
disabled={loading}
onClick={() => {
resetPages();
setRefresh((value) => value + 1);
}}
>
</button>
</div>
<form className="admin-panel admin-form" onSubmit={search}>
<div className="admin-filter-grid">
{(['userId', 'projectId'] as const).map((key) => (
<label key={key} className="admin-field">
<span>{fieldLabels[key]}</span>
<input
value={filters[key]}
onChange={(event) =>
setFilters({ ...filters, [key]: event.target.value })
}
/>
</label>
))}
<label className="admin-field">
<span></span>
<select
value={filters.eventName}
onChange={(event) =>
setFilters({ ...filters, eventName: event.target.value })
}
>
<option value=""></option>
{Object.entries(eventLabels).map(([value, label]) => (
<option key={value} value={value}>
{label}
</option>
))}
</select>
</label>
{(['startTime', 'endTime'] as const).map((key) => (
<label key={key} className="admin-field">
<span>
{key === 'startTime'
? '发生时间起点(含)'
: '发生时间终点(不含)'}
</span>
<input
type="datetime-local"
value={filters[key]}
onChange={(event) =>
setFilters({ ...filters, [key]: event.target.value })
}
/>
</label>
))}
</div>
<div className="admin-action-row">
<button
type="submit"
className="admin-primary-button"
disabled={loading}
>
</button>
</div>
{(['creativeTaskId', 'agentRunId', 'clientVersion'] as const)
.filter((key) => query[key])
.map((key) => (
<p key={key}>
{fieldLabels[key]}{query[key]}
</p>
))}
</form>
{error ? (
<p role="alert" className="admin-error-message">
{error}
</p>
) : null}
<div className="admin-panel">
<div className="admin-table-wrap">
<table className="admin-table">
<thead>
<tr>
{[
'入库时间',
'发生时间',
'用户',
'事件名称',
'项目',
'来源',
'结果',
'客户端版本',
'详情',
].map((label) => (
<th key={label}>{label}</th>
))}
</tr>
</thead>
<tbody>
{entries.map((entry) => (
<tr key={entry.eventId}>
<td>{formatTime(entry.receivedAt)}</td>
<td>{formatTime(entry.eventTime)}</td>
<td>
{entry.userId}
<AdminUserReferenceButton
token={token}
userId={entry.userId}
onUnauthorized={onUnauthorized}
/>
</td>
<td>{eventLabels[entry.eventName] ?? entry.eventName}</td>
<td>{entry.projectId ?? '—'}</td>
<td>{entry.source}</td>
<td>{entry.status ?? '—'}</td>
<td>{entry.clientVersion}</td>
<td>
<button
type="button"
className="admin-ghost-button"
onClick={() => setDetail(entry)}
>
</button>
</td>
</tr>
))}
</tbody>
</table>
</div>
{loading ? (
<p role="status"></p>
) : !error && entries.length === 0 ? (
<p></p>
) : null}
<div className="admin-action-row">
<button
type="button"
className="admin-secondary-button"
disabled={loading || page === 0}
onClick={() => setPage((value) => value - 1)}
>
</button>
<span> {page + 1} </span>
<button
type="button"
className="admin-secondary-button"
disabled={loading || !nextCursor}
onClick={() => {
if (!nextCursor) return;
setCursors([...cursors.slice(0, page + 1), nextCursor]);
setPage(page + 1);
}}
>
</button>
</div>
</div>
{detail ? (
<Modal
open
title="客户端埋点详情"
closeLabel="关闭详情"
onClose={() => setDetail(null)}
className="genarrative-ui"
>
<dl>
{(
Object.keys(fieldLabels) as Array<
keyof AdminAgcTrackingEventEntry
>
)
.filter((key) => key !== 'properties')
.map((key) => (
<div key={key}>
<dt>{fieldLabels[key]}</dt>
<dd style={{ overflowWrap: 'anywhere' }}>
{detail[key] == null ? '—' : String(detail[key])}
</dd>
</div>
))}
</dl>
<div className="admin-action-row">
{(['creativeTaskId', 'agentRunId', 'clientVersion'] as const).map(
(key) =>
detail[key] ? (
<button
type="button"
key={key}
className="admin-secondary-button"
onClick={() => related(key, detail[key]!)}
>
{fieldLabels[key]}
</button>
) : null,
)}
</div>
<h3></h3>
<pre style={{ whiteSpace: 'pre-wrap', overflowWrap: 'anywhere' }}>
{JSON.stringify(detail.properties, null, 2)}
</pre>
</Modal>
) : null}
</section>
);
}
@@ -2641,6 +2641,48 @@ const databaseTableLabelMap: Record<string, string> = {
profile_recharge_order: '充值订单', profile_recharge_order: '充值订单',
profile_feedback_submission: '反馈提交', profile_feedback_submission: '反馈提交',
profile_save_archive: '存档记录', profile_save_archive: '存档记录',
story_session: '剧情会话',
story_event: '剧情事件',
npc_state: 'NPC 状态',
inventory_slot: '背包槽位',
battle_state: '战斗状态',
treasure_record: '宝藏记录',
quest_record: '任务记录',
quest_log: '任务日志',
player_progression: '玩家进度',
chapter_progression: '章节进度',
custom_world_profile: '自定义世界档案',
custom_world_session: '自定义世界会话',
custom_world_agent_session: '自定义世界 Agent 会话',
custom_world_agent_message: '自定义世界 Agent 消息',
custom_world_agent_operation: '自定义世界 Agent 操作',
custom_world_draft_card: '自定义世界草稿卡片',
custom_world_gallery_entry: '自定义世界画廊条目',
puzzle_agent_session: '拼图 Agent 会话',
puzzle_agent_message: '拼图 Agent 消息',
puzzle_work_profile: '拼图作品档案',
puzzle_event: '拼图事件',
puzzle_runtime_run: '拼图运行记录',
puzzle_leaderboard_entry: '拼图排行榜条目',
match3d_agent_session: '抓大鹅 Agent 会话',
match3d_agent_message: '抓大鹅 Agent 消息',
match3d_work_profile: '抓大鹅作品档案',
match3d_runtime_run: '抓大鹅运行记录',
square_hole_agent_session: '方洞挑战 Agent 会话',
square_hole_agent_message: '方洞挑战 Agent 消息',
square_hole_work_profile: '方洞挑战作品档案',
square_hole_runtime_run: '方洞挑战运行记录',
visual_novel_agent_session: '视觉小说 Agent 会话',
visual_novel_agent_message: '视觉小说 Agent 消息',
visual_novel_work_profile: '视觉小说作品档案',
visual_novel_runtime_run: '视觉小说运行记录',
visual_novel_runtime_history_entry: '视觉小说历史条目',
visual_novel_runtime_event: '视觉小说运行事件',
big_fish_creation_session: '大鱼吃小鱼创建会话',
big_fish_agent_message: '大鱼吃小鱼 Agent 消息',
big_fish_asset_slot: '大鱼吃小鱼资产槽位',
big_fish_event: '大鱼吃小鱼事件',
big_fish_runtime_run: '大鱼吃小鱼运行记录',
asset_object: '资产对象', asset_object: '资产对象',
asset_entity_binding: '资产实体绑定', asset_entity_binding: '资产实体绑定',
asset_event: '资产事件', asset_event: '资产事件',
@@ -2682,6 +2724,48 @@ const databaseTableDescriptionMap: Record<string, string> = {
profile_recharge_order: '充值订单表', profile_recharge_order: '充值订单表',
profile_feedback_submission: '反馈提交记录表', profile_feedback_submission: '反馈提交记录表',
profile_save_archive: '用户存档记录表', profile_save_archive: '用户存档记录表',
story_session: '剧情会话表',
story_event: '剧情事件表',
npc_state: 'NPC 状态表',
inventory_slot: '背包槽位表',
battle_state: '战斗状态表',
treasure_record: '宝藏记录表',
quest_record: '任务记录表',
quest_log: '任务日志表',
player_progression: '玩家进度表',
chapter_progression: '章节进度表',
custom_world_profile: '自定义世界档案表',
custom_world_session: '自定义世界会话表',
custom_world_agent_session: '自定义世界 Agent 会话表',
custom_world_agent_message: '自定义世界 Agent 消息表',
custom_world_agent_operation: '自定义世界 Agent 操作表',
custom_world_draft_card: '自定义世界草稿卡片表',
custom_world_gallery_entry: '自定义世界画廊条目表',
puzzle_agent_session: '拼图 Agent 会话表',
puzzle_agent_message: '拼图 Agent 消息表',
puzzle_work_profile: '拼图作品档案表',
puzzle_event: '拼图事件表',
puzzle_runtime_run: '拼图运行记录表',
puzzle_leaderboard_entry: '拼图排行榜条目表',
match3d_agent_session: '抓大鹅 Agent 会话表',
match3d_agent_message: '抓大鹅 Agent 消息表',
match3d_work_profile: '抓大鹅作品档案表',
match3d_runtime_run: '抓大鹅运行记录表',
square_hole_agent_session: '方洞挑战 Agent 会话表',
square_hole_agent_message: '方洞挑战 Agent 消息表',
square_hole_work_profile: '方洞挑战作品档案表',
square_hole_runtime_run: '方洞挑战运行记录表',
visual_novel_agent_session: '视觉小说 Agent 会话表',
visual_novel_agent_message: '视觉小说 Agent 消息表',
visual_novel_work_profile: '视觉小说作品档案表',
visual_novel_runtime_run: '视觉小说运行记录表',
visual_novel_runtime_history_entry: '视觉小说历史条目表',
visual_novel_runtime_event: '视觉小说运行事件表',
big_fish_creation_session: '大鱼吃小鱼创建会话表',
big_fish_agent_message: '大鱼吃小鱼 Agent 消息表',
big_fish_asset_slot: '大鱼吃小鱼资产槽位表',
big_fish_event: '大鱼吃小鱼事件表',
big_fish_runtime_run: '大鱼吃小鱼运行记录表',
asset_object: '资产对象表', asset_object: '资产对象表',
asset_entity_binding: '资产实体绑定表', asset_entity_binding: '资产实体绑定表',
asset_event: '资产事件表', asset_event: '资产事件表',
@@ -175,107 +175,6 @@ test('灰度发布页可选择模板库并默认启用零比例灰度', async ()
); );
}); });
test('灰度发布页可选择游戏发布开关,默认保持「未开启即开放」语义', async () => {
const user = userEvent.setup();
render(
<AdminGrayReleaseConfigPage token="admin-token" onUnauthorized={vi.fn()} />,
);
await screen.findByRole('button', { name: 'editor.new-toolbar' });
await user.selectOptions(screen.getByLabelText('Gate Key 前缀'), [
'game-distribution',
]);
expect((screen.getByLabelText('Gate Key') as HTMLInputElement).value).toBe(
'game-distribution:publish',
);
expect(
(screen.getByLabelText('Gate Key 目标') as HTMLSelectElement).value,
).toBe('publish');
// 该开关的语义是「未配置/关闭 = 默认开放」,所以选中后不能默认打开收紧。
expect((screen.getByLabelText('启用') as HTMLInputElement).checked).toBe(
false,
);
expect((screen.getByLabelText('灰度比例') as HTMLInputElement).value).toBe(
'0',
);
expect(
(screen.getByLabelText('描述') as HTMLTextAreaElement).value,
).toContain('游戏发布入口灰度');
});
test('灰度发布页保存游戏发布开关时写入白名单与比例', async () => {
const user = userEvent.setup();
vi.mocked(upsertAdminFeatureGateConfig).mockResolvedValueOnce({
gates: [
...configResponse.gates,
{
gateKey: 'game-distribution:publish',
enabled: true,
rolloutPercent: 20,
allowUserIds: ['user-internal'],
allowUserTags: [],
denyUserIds: [],
description: '游戏发布入口灰度',
updatedAt: '2026-09-22T10:00:00Z',
},
],
});
render(
<AdminGrayReleaseConfigPage token="admin-token" onUnauthorized={vi.fn()} />,
);
await screen.findByRole('button', { name: 'editor.new-toolbar' });
await user.selectOptions(screen.getByLabelText('Gate Key 前缀'), [
'game-distribution',
]);
fireEvent.click(screen.getByLabelText('启用'));
fireEvent.change(screen.getByLabelText('灰度比例'), {
target: { value: '20' },
});
fireEvent.change(screen.getByLabelText('允许用户 ID'), {
target: { value: 'user-internal' },
});
fireEvent.change(screen.getByLabelText('描述'), {
target: { value: '游戏发布入口灰度' },
});
await user.click(screen.getByRole('button', { name: '保存配置' }));
await user.click(screen.getByRole('button', { name: '确认' }));
await waitFor(() =>
expect(upsertAdminFeatureGateConfig).toHaveBeenCalledWith('admin-token', {
gateKey: 'game-distribution:publish',
enabled: true,
rolloutPercent: 20,
allowUserIds: ['user-internal'],
allowUserTags: [],
denyUserIds: [],
description: '游戏发布入口灰度',
}),
);
});
test('未创建的预设开关在后台可见并可一键配置', async () => {
const user = userEvent.setup();
render(
<AdminGrayReleaseConfigPage token="admin-token" onUnauthorized={vi.fn()} />,
);
const row = await screen.findByText('game-distribution:publish');
expect(row).not.toBeNull();
// 该开关默认未创建:列表里给出「配置」入口,点击后按默认关闭填充表单。
const configureButton = row.closest('tr')?.querySelector('button');
expect(configureButton).not.toBeNull();
await user.click(configureButton!);
expect((screen.getByLabelText('Gate Key') as HTMLInputElement).value).toBe(
'game-distribution:publish',
);
expect((screen.getByLabelText('启用') as HTMLInputElement).checked).toBe(
false,
);
});
test('灰度发布页保存时转换数组和百分比', async () => { test('灰度发布页保存时转换数组和百分比', async () => {
const user = userEvent.setup(); const user = userEvent.setup();
vi.mocked(upsertAdminFeatureGateConfig).mockResolvedValueOnce({ vi.mocked(upsertAdminFeatureGateConfig).mockResolvedValueOnce({
@@ -28,7 +28,6 @@ interface GateTargetOption {
const GATE_PREFIX_LABELS: Record<string, string> = { const GATE_PREFIX_LABELS: Record<string, string> = {
'image-editor': '画布', 'image-editor': '画布',
agc: '客户端', agc: '客户端',
'game-distribution': '游戏分发',
}; };
const FIXED_GATE_TARGETS: GateTargetOption[] = [ const FIXED_GATE_TARGETS: GateTargetOption[] = [
@@ -46,14 +45,6 @@ const FIXED_GATE_TARGETS: GateTargetOption[] = [
label: 'Agent 侧边栏', label: 'Agent 侧边栏',
description: '画布 Agent 入口灰度', description: '画布 Agent 入口灰度',
}, },
{
prefix: 'game-distribution',
suffix: 'publish',
key: 'game-distribution:publish',
label: '游戏发布',
description:
'游戏发布入口灰度:未配置或关闭时对已登录作者默认开放,开启后只放行白名单 / 灰度命中',
},
]; ];
export function AdminGrayReleaseConfigPage({ export function AdminGrayReleaseConfigPage({
@@ -206,11 +197,6 @@ export function AdminGrayReleaseConfigPage({
setErrorMessage(''); setErrorMessage('');
} }
// 预设里尚未创建行的开关也要可见:运营需要先看到 key 才能配置灰度。
const unconfiguredGateTargets = FIXED_GATE_TARGETS.filter(
(option) => !gates.some((gate) => gate.gateKey === option.key),
);
function buildPayload(): AdminUpsertFeatureGateConfigRequest { function buildPayload(): AdminUpsertFeatureGateConfigRequest {
return { return {
gateKey: gateKey.trim(), gateKey: gateKey.trim(),
@@ -457,53 +443,6 @@ export function AdminGrayReleaseConfigPage({
</div> </div>
)} )}
</section> </section>
<section className="admin-panel">
<div className="admin-panel-heading">
<h3></h3>
<span>{unconfiguredGateTargets.length}</span>
</div>
{unconfiguredGateTargets.length ? (
<div className="admin-table-wrap">
<table className="admin-table admin-table-compact">
<thead>
<tr>
<th>Gate</th>
<th></th>
<th></th>
</tr>
</thead>
<tbody>
{unconfiguredGateTargets.map((option) => (
<tr key={option.key}>
<td>
{option.key}
<small>
{GATE_PREFIX_LABELS[option.prefix] ?? option.prefix} ·{' '}
{option.label}
</small>
</td>
<td>{option.description}</td>
<td>
<button
className="admin-text-button"
type="button"
onClick={() => applyGateTarget(option)}
>
</button>
</td>
</tr>
))}
</tbody>
</table>
</div>
) : (
<div className="admin-empty-state">
{isLoading ? '加载中' : '预设开关都已创建'}
</div>
)}
</section>
</div> </div>
{confirmDialog} {confirmDialog}
@@ -154,6 +154,7 @@ const allowedUncalledTauriCommands = [
'steer_game_creator_agent_runtime_task', 'steer_game_creator_agent_runtime_task',
'write_local_agent_memory', 'write_local_agent_memory',
'write_local_game_memory', 'write_local_game_memory',
'write_local_project_file',
// Agent 运行时会话 / 目标 / 协作命令由 native 侧与 CLI swarm 驱动,前端没有调用方。 // Agent 运行时会话 / 目标 / 协作命令由 native 侧与 CLI swarm 驱动,前端没有调用方。
'archive_game_creator_agent_session', 'archive_game_creator_agent_session',
'clear_game_creator_agent_goal', 'clear_game_creator_agent_goal',
@@ -516,7 +517,7 @@ function parseTauriHandlerCommandNames(source) {
throw new Error('AI game creator shell Tauri handler list is missing'); throw new Error('AI game creator shell Tauri handler list is missing');
} }
return Array.from( return Array.from(
match[1].matchAll(/\b(?:[a-z][a-z0-9_]*::)*([a-z][a-z0-9_]*)\b/g), match[1].matchAll(/\b([a-z][a-z0-9_]+)\b/g),
([, command]) => command, ([, command]) => command,
); );
} }
@@ -549,12 +550,6 @@ function assertCommandNamesDisjoint(label, leftNames, rightNames) {
} }
function runAppInvokeParserRegressionChecks() { function runAppInvokeParserRegressionChecks() {
assert.deepEqual(
parseTauriHandlerCommandNames(
'tauri::generate_handler![plain_command, analytics::gui::capture_analytics_context,]',
),
['plain_command', 'capture_analytics_context'],
);
assert.deepEqual( assert.deepEqual(
parseAppInvokeCommandNames(` parseAppInvokeCommandNames(`
invoke('direct_command', {}); invoke('direct_command', {});
@@ -34,6 +34,7 @@ mod direct_thread_wire;
mod direct_tool_bridge; mod direct_tool_bridge;
mod direct_tool_calls; mod direct_tool_calls;
mod direct_tools_mcp; mod direct_tools_mcp;
mod direct_turn_failure;
mod direct_turn_metrics; mod direct_turn_metrics;
mod direct_turn_stream; mod direct_turn_stream;
mod direct_validation; mod direct_validation;
@@ -72,6 +73,7 @@ pub(crate) use direct_thread_wire::*;
pub(crate) use direct_tool_bridge::*; pub(crate) use direct_tool_bridge::*;
pub(crate) use direct_tool_calls::*; pub(crate) use direct_tool_calls::*;
pub(crate) use direct_tools_mcp::*; pub(crate) use direct_tools_mcp::*;
pub(crate) use direct_turn_failure::*;
pub(crate) use direct_turn_metrics::*; pub(crate) use direct_turn_metrics::*;
pub(crate) use direct_turn_stream::*; pub(crate) use direct_turn_stream::*;
pub(crate) use direct_validation::DirectValidationConfig; pub(crate) use direct_validation::DirectValidationConfig;
@@ -3548,12 +3548,20 @@ impl CodexAppServerConnection {
let direct_turn_user_item_id = direct_persisted_user_item let direct_turn_user_item_id = direct_persisted_user_item
.as_ref() .as_ref()
.and_then(direct_thread_item_identity); .and_then(direct_thread_item_identity);
// 回合终态兜底:`turn.started` 进队列之后就武装,写完终态即解除。宿主在这两者之间任何
// 提前收场(panic、future 被丢弃、以后新增的早退)都由它补一条失败终态,否则前端只能
// 永远停在"还在跑"。
let mut direct_turn_failure_guard: Option<DirectTurnFailureGuard> = None;
if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject {
append_direct_thread_event( append_direct_thread_event(
&direct_thread_id, &direct_thread_id,
DirectThreadEvent::turn_started(direct_turn_started_at_ms) DirectThreadEvent::turn_started(direct_turn_started_at_ms)
.with_user_item_id(direct_turn_user_item_id.as_deref()), .with_user_item_id(direct_turn_user_item_id.as_deref()),
); );
direct_turn_failure_guard = Some(DirectTurnFailureGuard::arm(
direct_thread_id.clone(),
direct_turn_user_item_id.clone(),
));
if let Some(user_item) = direct_persisted_user_item.as_ref() { if let Some(user_item) = direct_persisted_user_item.as_ref() {
if let Some(entry_item) = direct_thread_event_item(history_root, user_item) { if let Some(entry_item) = direct_thread_event_item(history_root, user_item) {
// 这里的条目时间可能是启动应答后的观测时间;前端按同一用户条目身份 // 这里的条目时间可能是启动应答后的观测时间;前端按同一用户条目身份
@@ -4061,11 +4069,30 @@ impl CodexAppServerConnection {
.map(|(_, at)| *at) .map(|(_, at)| *at)
.unwrap_or_else(direct_tool_call_now_ms) .unwrap_or_else(direct_tool_call_now_ms)
}; };
append_direct_thread_event( // 终态只有 `turn.completed` 一种事件:失败时同一个事件带 `failure` 载荷(原因由宿主
// 脱敏 + 截断后写进去),其余(`completed` / `interrupted` / `aborted`)不带载荷。
// 失败不再只写一个 `status="failed"`:那让失败与正常结束在协议上长得一样,前端只能
// 另开一条通道(命令返回 / 另一条 IPC)去拿原因,也就等于承认事件流讲不清一轮怎么结束。
let failure = direct_turn_failure(
&status,
collect_result.as_ref().map(String::as_str),
history_root,
);
match failure {
Some(failure) => append_direct_thread_event(
&direct_thread_id,
DirectThreadEvent::turn_completed_failed(failure, completed_at)
.with_user_item_id(direct_turn_user_item_id.as_deref()),
),
None => append_direct_thread_event(
&direct_thread_id, &direct_thread_id,
DirectThreadEvent::turn_completed(status, completed_at) DirectThreadEvent::turn_completed(status, completed_at)
.with_user_item_id(direct_turn_user_item_id.as_deref()), .with_user_item_id(direct_turn_user_item_id.as_deref()),
); ),
};
if let Some(guard) = direct_turn_failure_guard.as_mut() {
guard.disarm();
}
} }
let text = collect_result?; let text = collect_result?;
guard.armed = false; guard.armed = false;
File diff suppressed because it is too large Load Diff
@@ -93,8 +93,6 @@ pub(super) struct ExecutionLedger {
pub(super) plan: Option<Value>, pub(super) plan: Option<Value>,
#[serde(default)] #[serde(default)]
pub(super) last_failed_write_revision: Option<u64>, pub(super) last_failed_write_revision: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(super) analytics_run: Option<crate::analytics::run::Metadata>,
} }
struct SessionData { struct SessionData {
@@ -108,17 +106,6 @@ struct SessionData {
elapsed_offset_ms: u64, elapsed_offset_ms: u64,
} }
// 与业务持久化锁分离;只保存内存状态,持锁期间不执行 I/O 或投递事件。
struct SessionAnalytics {
project_id: String,
route: Option<crate::analytics::contract::Route>,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
output_revision: Option<u64>,
}
#[derive(Clone, PartialEq, Eq)] #[derive(Clone, PartialEq, Eq)]
struct CodexExecutorIdentity { struct CodexExecutorIdentity {
path: PathBuf, path: PathBuf,
@@ -166,12 +153,9 @@ fn executor_digest(path: &Path) -> Result<String, String> {
pub(super) struct ExecutionSession { pub(super) struct ExecutionSession {
pub(super) root: PathBuf, pub(super) root: PathBuf,
/// 本次确实新建执行账本;恢复和旧预算迁移均不构成新的用户受理。
pub(super) newly_accepted: bool,
state_path: PathBuf, state_path: PathBuf,
_owner: File, _owner: File,
data: Mutex<SessionData>, data: Mutex<SessionData>,
analytics: Mutex<SessionAnalytics>,
changed: tokio::sync::watch::Sender<u64>, changed: tokio::sync::watch::Sender<u64>,
cancellation: Arc<std::sync::atomic::AtomicBool>, cancellation: Arc<std::sync::atomic::AtomicBool>,
abort_requested: std::sync::atomic::AtomicBool, abort_requested: std::sync::atomic::AtomicBool,
@@ -232,16 +216,6 @@ pub(crate) struct WritePermit {
id: String, id: String,
} }
impl WritePermit { impl WritePermit {
pub(super) fn record_analytics_revision(
&self,
revision: u64,
change_kind: crate::analytics::contract::ChangeKind,
files_changed_count: u64,
) {
self.session
.record_analytics_revision(revision, change_kind, files_changed_count);
}
pub(crate) fn run<T>(&self, write: impl FnOnce() -> Result<T, String>) -> Result<T, String> { pub(crate) fn run<T>(&self, write: impl FnOnce() -> Result<T, String>) -> Result<T, String> {
let mut data = self.session.lock()?; let mut data = self.session.lock()?;
self.session.tick_locked(&mut data)?; self.session.tick_locked(&mut data)?;
@@ -427,7 +401,6 @@ pub(super) async fn begin(
prompt: &str, prompt: &str,
requires_contract: bool, requires_contract: bool,
config: DirectValidationConfig, config: DirectValidationConfig,
analytics_run: Option<crate::analytics::run::Metadata>,
) -> Result<ExecutionSessionGuard, String> { ) -> Result<ExecutionSessionGuard, String> {
let root = root.to_path_buf(); let root = root.to_path_buf();
let prompt_hash = hash(prompt.as_bytes()); let prompt_hash = hash(prompt.as_bytes());
@@ -435,14 +408,13 @@ pub(super) async fn begin(
let turn = super::direct_taonier_active_invocation_id_at(&root)?; let turn = super::direct_taonier_active_invocation_id_at(&root)?;
let host = crate::game_creator_runtime_config_dir() let host = crate::game_creator_runtime_config_dir()
.ok_or("direct-execution-host: 需要客户端私有配置目录,CLI 请提供 --config-dir")?; .ok_or("direct-execution-host: 需要客户端私有配置目录,CLI 请提供 --config-dir")?;
open_with_analytics_at( open_at(
&host.join("direct-executions"), &host.join("direct-executions"),
&root, &root,
&turn, &turn,
&prompt_hash, &prompt_hash,
requires_contract, requires_contract,
&config, &config,
analytics_run,
) )
}) })
.await .await
@@ -490,26 +462,6 @@ pub(super) fn open_at(
request_hash: &str, request_hash: &str,
requires_contract: bool, requires_contract: bool,
config: &DirectValidationConfig, config: &DirectValidationConfig,
) -> Result<Arc<ExecutionSession>, String> {
open_with_analytics_at(
host,
root,
turn,
request_hash,
requires_contract,
config,
None,
)
}
pub(super) fn open_with_analytics_at(
host: &Path,
root: &Path,
turn: &str,
request_hash: &str,
requires_contract: bool,
config: &DirectValidationConfig,
analytics_run: Option<crate::analytics::run::Metadata>,
) -> Result<Arc<ExecutionSession>, String> { ) -> Result<Arc<ExecutionSession>, String> {
config.validate()?; config.validate()?;
let root = root let root = root
@@ -569,7 +521,6 @@ pub(super) fn open_with_analytics_at(
}; };
let project_id = super::read_existing_manifest_for_project(&root)?.project_id; let project_id = super::read_existing_manifest_for_project(&root)?.project_id;
let is_new = existing.is_none(); let is_new = existing.is_none();
let mut newly_accepted = is_new;
let mut ledger = existing.unwrap_or_else(|| ExecutionLedger { let mut ledger = existing.unwrap_or_else(|| ExecutionLedger {
schema_version: SCHEMA.into(), schema_version: SCHEMA.into(),
client_turn_id: turn.into(), client_turn_id: turn.into(),
@@ -595,7 +546,6 @@ pub(super) fn open_with_analytics_at(
delivery_reviews: 0, delivery_reviews: 0,
plan: None, plan: None,
last_failed_write_revision: None, last_failed_write_revision: None,
analytics_run,
}); });
if is_new { if is_new {
// 只继承旧项目账本的消费量,绝不把可编辑的旧成功回执提升为宿主证据。 // 只继承旧项目账本的消费量,绝不把可编辑的旧成功回执提升为宿主证据。
@@ -608,8 +558,6 @@ pub(super) fn open_with_analytics_at(
512 * 1024, 512 * 1024,
)?; )?;
if let Some(legacy) = legacy { if let Some(legacy) = legacy {
newly_accepted = false;
ledger.analytics_run = None;
let used = legacy["usedRuns"] let used = legacy["usedRuns"]
.as_u64() .as_u64()
.and_then(|n| u32::try_from(n).ok()); .and_then(|n| u32::try_from(n).ok());
@@ -666,18 +614,8 @@ pub(super) fn open_with_analytics_at(
let (changed, _) = tokio::sync::watch::channel(ledger.revision); let (changed, _) = tokio::sync::watch::channel(ledger.revision);
let session = Arc::new(ExecutionSession { let session = Arc::new(ExecutionSession {
root, root,
newly_accepted,
state_path, state_path,
_owner: owner, _owner: owner,
analytics: Mutex::new(SessionAnalytics {
project_id: ledger.project_id.clone(),
route: ledger
.analytics_run
.as_ref()
.map(|run| run.context.route.clone()),
capture: None,
output_revision: None,
}),
data: Mutex::new(SessionData { data: Mutex::new(SessionData {
ledger, ledger,
started: Instant::now(), started: Instant::now(),
@@ -800,10 +738,6 @@ impl ExecutionSession {
pub(super) fn cancel_flag(&self) -> Arc<std::sync::atomic::AtomicBool> { pub(super) fn cancel_flag(&self) -> Arc<std::sync::atomic::AtomicBool> {
Arc::clone(&self.cancellation) Arc::clone(&self.cancellation)
} }
pub(super) fn was_aborted(&self) -> bool {
self.abort_requested
.load(std::sync::atomic::Ordering::Acquire)
}
pub(super) fn record_delivery_review(&self) -> Result<u32, String> { pub(super) fn record_delivery_review(&self) -> Result<u32, String> {
let mut data = self.lock()?; let mut data = self.lock()?;
if data.ledger.phase.is_terminal() { if data.ledger.phase.is_terminal() {
@@ -831,73 +765,6 @@ impl ExecutionSession {
self.commit(&mut data, next)?; self.commit(&mut data, next)?;
Ok(json!({"plan":plan,"revision":data.ledger.revision,"acceptancePassed":false})) Ok(json!({"plan":plan,"revision":data.ledger.revision,"acceptancePassed":false}))
} }
pub(super) fn set_analytics_capture(
&self,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
) {
let Ok(mut analytics) = self.analytics.lock() else {
return;
};
analytics.capture = capture.and_then(|(mut context, writer)| {
// 恢复或账号切换后仍归属于真实受理的原 run。
context.route = analytics.route.clone()?;
Some((context, writer))
});
}
pub(super) fn analytics_capture(
&self,
) -> Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)> {
self.analytics.lock().ok()?.capture.clone()
}
pub(super) fn record_analytics_revision(
&self,
revision: u64,
change_kind: crate::analytics::contract::ChangeKind,
files_changed_count: u64,
) {
use crate::analytics::contract::{RevisionCreated, RevisionSource, Source};
if files_changed_count == 0 {
return;
}
let Ok(mut analytics) = self.analytics.lock() else {
return;
};
if analytics.route.is_none() {
return;
}
analytics.output_revision = Some(analytics.output_revision.unwrap_or(0).max(revision));
let capture = analytics.capture.clone();
let project_id = analytics.project_id.clone();
drop(analytics);
crate::analytics::project::revision(
capture,
&project_id,
Source::Direct,
RevisionCreated {
revision_id: revision.to_string(),
revision_source: RevisionSource::Agent,
change_kind,
files_changed_count: Some(files_changed_count),
},
);
}
pub(super) fn analytics_output_revision(&self) -> Option<String> {
self.analytics
.lock()
.ok()?
.output_revision
.map(|revision| revision.to_string())
}
pub(super) fn snapshot(&self) -> Result<ExecutionLedger, String> { pub(super) fn snapshot(&self) -> Result<ExecutionLedger, String> {
let data = self.lock()?; let data = self.lock()?;
let mut state = data.ledger.clone(); let mut state = data.ledger.clone();
@@ -1,192 +1,5 @@
use super::*; use super::*;
fn analytics_metadata(user: &str) -> crate::analytics::run::Metadata {
use crate::analytics::contract::{Context, Route, RunSource, Source};
crate::analytics::run::Metadata::new(
Context {
route: Route::from_identity(Some(user.into()), Some("https://example.com")),
editor_session_id: uuid::Uuid::new_v4().to_string(),
client_version: "1.0.0".into(),
},
Source::Direct,
RunSource::UserSubmit,
)
}
#[test]
fn analytics_survives_business_lock_contention_and_preserves_replayed_run_identity() {
use crate::analytics::{contract::ChangeKind, store::AnalyticsWriter};
let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("project");
let host = temp.path().join("host");
crate::init_local_game_project_at(&root, "analytics-lock", "采集锁隔离").unwrap();
let original = open_with_analytics_at(
&host,
&root,
"turn",
&hash(b"request"),
false,
&Default::default(),
Some(analytics_metadata("A")),
)
.unwrap();
drop(original);
let current = analytics_metadata("B");
let session = open_with_analytics_at(
&host,
&root,
"turn",
&hash(b"request"),
false,
&Default::default(),
Some(current.clone()),
)
.unwrap();
let config = temp.path().join("config");
std::fs::create_dir_all(&config).unwrap();
let writer = AnalyticsWriter::start(config.clone(), current.context.editor_session_id.clone());
let capture = (current.context.clone(), writer.clone());
// 模拟业务提交长期占锁;采集必须在释放该锁之前完成。
let business_lock = session.data.lock().unwrap();
let task_session = session.clone();
let (sender, receiver) = std::sync::mpsc::channel();
let worker = std::thread::spawn(move || {
task_session.set_analytics_capture(Some(capture));
task_session.record_analytics_revision(7, ChangeKind::Code, 1);
task_session.record_analytics_revision(5, ChangeKind::Code, 1);
task_session.record_analytics_revision(99, ChangeKind::Code, 0);
sender
.send((
task_session.analytics_capture(),
task_session.analytics_output_revision(),
))
.unwrap();
});
let result = receiver.recv_timeout(std::time::Duration::from_secs(5));
// 即使回归成等待业务锁,也先释放锁和回收线程,让测试明确失败而非挂死。
drop(business_lock);
worker.join().unwrap();
let (capture, revision) = result.expect("采集不得等待业务持久化锁");
let (context, _) = capture.expect("锁竞争不得丢失采集身份");
assert_eq!(context.route.user_id.as_deref(), Some("A"));
assert_eq!(context.editor_session_id, current.context.editor_session_id);
assert_eq!(revision.as_deref(), Some("7"));
assert!(writer.flush());
let batches = config
.join("analytics/instances")
.join(&current.context.editor_session_id)
.join("batches");
let deadline = Instant::now() + std::time::Duration::from_secs(5);
loop {
let events: Vec<Value> = std::fs::read_dir(&batches)
.into_iter()
.flatten()
.flatten()
.filter(|entry| !entry.file_name().to_string_lossy().starts_with('.'))
.filter_map(|entry| std::fs::read_to_string(entry.path().join("events.jsonl")).ok())
.flat_map(|text| {
text.lines()
.map(|line| serde_json::from_str::<Value>(line).unwrap())
.collect::<Vec<_>>()
})
.collect();
if events.len() == 2 {
for (event, expected_revision) in events.iter().zip(["7", "5"]) {
assert_eq!(event["event_name"], "project_revision_created");
assert_eq!(event["user_id"], "A");
assert_eq!(event["project_id"], "analytics-lock");
assert_eq!(event["properties"]["revision_id"], expected_revision);
}
break;
}
assert!(Instant::now() < deadline, "成果事件未落盘");
std::thread::sleep(std::time::Duration::from_millis(5));
}
drop(session);
let resumed = open_with_analytics_at(
&host,
&root,
"turn",
&hash(b"request"),
false,
&Default::default(),
Some(current),
)
.unwrap();
assert_eq!(
resumed.analytics_output_revision(),
None,
"恢复不补造历史成果编号"
);
}
#[test]
fn run_metadata_is_persisted_with_new_ledger_and_replay_keeps_original_identity() {
let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("project");
crate::init_local_game_project_at(&root, "analytics-run", "执行身份").unwrap();
let original = analytics_metadata("A");
let host = temp.path().join("host");
let session = open_with_analytics_at(
&host,
&root,
"turn",
&hash(b"request"),
false,
&Default::default(),
Some(original.clone()),
)
.unwrap();
assert!(session.newly_accepted);
assert_eq!(
session.snapshot().unwrap().analytics_run,
Some(original.clone())
);
drop(session);
let replay = open_with_analytics_at(
&host,
&root,
"turn",
&hash(b"request"),
false,
&Default::default(),
Some(analytics_metadata("B")),
)
.unwrap();
assert!(!replay.newly_accepted);
assert_eq!(replay.snapshot().unwrap().analytics_run, Some(original));
}
#[test]
fn legacy_run_without_metadata_is_not_assigned_current_users_identity() {
let (temp, session) = fixture(Default::default());
let root = session.root.clone();
assert!(session.snapshot().unwrap().analytics_run.is_none());
drop(session);
let replay = open_with_analytics_at(
&temp.path().join("host"),
&root,
"turn-test",
&hash(b"request"),
false,
&Default::default(),
Some(analytics_metadata("B")),
)
.unwrap();
assert!(replay.snapshot().unwrap().analytics_run.is_none());
let config = temp.path().join("config");
std::fs::create_dir_all(&config).unwrap();
let current = analytics_metadata("B");
let writer = crate::analytics::store::AnalyticsWriter::start(
config,
current.context.editor_session_id.clone(),
);
replay.set_analytics_capture(Some((current.context, writer)));
replay.record_analytics_revision(1, crate::analytics::contract::ChangeKind::Code, 1);
assert!(replay.analytics_capture().is_none());
assert!(replay.analytics_output_revision().is_none());
}
fn fixture(config: DirectValidationConfig) -> (tempfile::TempDir, Arc<ExecutionSession>) { fn fixture(config: DirectValidationConfig) -> (tempfile::TempDir, Arc<ExecutionSession>) {
let temp = tempfile::tempdir().unwrap(); let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("project"); let root = temp.path().join("project");
@@ -200,7 +13,6 @@ fn fixture(config: DirectValidationConfig) -> (tempfile::TempDir, Arc<ExecutionS
&config, &config,
) )
.unwrap(); .unwrap();
assert!(session.newly_accepted);
session session
.freeze_contract(json!({"requirements":[{"id":"test"}]})) .freeze_contract(json!({"requirements":[{"id":"test"}]}))
.unwrap(); .unwrap();
@@ -347,7 +159,6 @@ fn reopened_budget_and_deadline_cannot_be_increased_by_configuration() {
) )
.unwrap(); .unwrap();
let state = reopened.snapshot().unwrap(); let state = reopened.snapshot().unwrap();
assert!(!reopened.newly_accepted);
assert_eq!(state.delivery_reviews, 1); assert_eq!(state.delivery_reviews, 1);
assert_eq!( assert_eq!(
( (
@@ -467,7 +278,6 @@ fn legacy_budget_is_inherited_without_trusting_project_success_evidence() {
&Default::default(), &Default::default(),
) )
.unwrap(); .unwrap();
assert!(!session.newly_accepted);
assert!(session.admit(EffectKind::Execute, None).is_err()); assert!(session.admit(EffectKind::Execute, None).is_err());
let state = session.snapshot().unwrap(); let state = session.snapshot().unwrap();
assert_eq!(state.used_passes, 2); assert_eq!(state.used_passes, 2);
@@ -229,40 +229,6 @@ fn target_fingerprints(targets: &[PatchTarget]) -> BTreeMap<String, Option<Strin
.collect() .collect()
} }
fn analytics_patch_changes(
before: &BTreeMap<String, Option<String>>,
after: &BTreeMap<String, Option<String>>,
) -> Option<(crate::analytics::contract::ChangeKind, u64)> {
use crate::analytics::contract::ChangeKind;
let mut result = None;
for (path, before) in before {
let Some((before, after)) = before
.as_ref()
.zip(after.get(path).and_then(Option::as_ref))
else {
continue;
};
if before == after {
continue;
}
let Some(kind) = crate::analytics::project::file_change_kind(path) else {
continue;
};
result = Some(match result {
None => (kind, 1),
Some((current, count)) => (
if current == kind {
current
} else {
ChangeKind::Mixed
},
count + 1,
),
});
}
result
}
fn run_transaction( fn run_transaction(
root: &Path, root: &Path,
parsed: codex_patch_parser::ApplyPatchArgs, parsed: codex_patch_parser::ApplyPatchArgs,
@@ -353,13 +319,6 @@ fn run_transaction(
None None
}; };
lease.finish(passed && !uncertain, changed, None)?; lease.finish(passed && !uncertain, changed, None)?;
if passed && !uncertain {
if let (Some(revision), Some((kind, count))) =
(revision, analytics_patch_changes(&before, &after))
{
session.record_analytics_revision(revision, kind, count);
}
}
Ok(json!({ Ok(json!({
"status": if passed && !uncertain { "completed" } else { "failed" }, "status": if passed && !uncertain { "completed" } else { "failed" },
"changedPaths": if started { changed_paths } else { BTreeSet::new() }, "changedPaths": if started { changed_paths } else { BTreeSet::new() },
@@ -392,47 +351,6 @@ pub(super) async fn apply(root: &Path, arguments: &Value) -> Result<Value, Strin
mod tests { mod tests {
use super::*; use super::*;
#[test]
fn analytics_patch_counts_only_known_changed_outputs() {
let fingerprints = |entries: &[(&str, Option<&str>)]| {
entries
.iter()
.map(|(path, value)| (path.to_string(), value.map(str::to_string)))
.collect()
};
let before = fingerprints(&[
("game/a.js", Some("a")),
("game/same.js", Some("same")),
("game/unknown.js", None),
("game/unreadable.js", Some("old")),
(".agent/state.json", Some("old")),
]);
let after = fingerprints(&[
("game/a.js", Some("b")),
("game/same.js", Some("same")),
("game/unknown.js", Some("new")),
("game/unreadable.js", None),
(".agent/state.json", Some("new")),
]);
assert_eq!(
analytics_patch_changes(&before, &after),
Some((crate::analytics::contract::ChangeKind::Code, 1))
);
assert_eq!(analytics_patch_changes(&before, &before), None);
let before = fingerprints(&[
("game/a.js", Some("missing")),
("assets/a.png", Some("old")),
]);
let after = fingerprints(&[
("game/a.js", Some("new")),
("assets/a.png", Some("missing")),
]);
assert_eq!(
analytics_patch_changes(&before, &after),
Some((crate::analytics::contract::ChangeKind::Mixed, 2))
);
}
fn project() -> (tempfile::TempDir, PathBuf) { fn project() -> (tempfile::TempDir, PathBuf) {
let temp = tempfile::tempdir().unwrap(); let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("project"); let root = temp.path().join("project");
@@ -497,20 +415,15 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn bundled_patch_roundtrip_preserves_partial_failure_and_rejects_closed_turn() { async fn bundled_patch_roundtrip_preserves_partial_failure_and_rejects_closed_turn() {
let (temp, root) = project(); let (temp, root) = project();
let config = temp.path().join("analytics-config"); let session = direct_execution::open_at(
let (metadata, context, writer) =
super::super::direct_tool_bridge::analytics_test_writer(&config);
let session = direct_execution::open_with_analytics_at(
&temp.path().join("host"), &temp.path().join("host"),
&root, &root,
"patch-roundtrip", "patch-roundtrip",
&format!("{:x}", Sha256::digest(b"request")), &format!("{:x}", Sha256::digest(b"request")),
false, false,
&direct_validation::DirectValidationConfig::default(), &direct_validation::DirectValidationConfig::default(),
Some(metadata),
) )
.unwrap(); .unwrap();
session.set_analytics_capture(Some((context.clone(), writer.clone())));
session session
.freeze_contract(json!({"fixture":"patch protocol only"})) .freeze_contract(json!({"fixture":"patch protocol only"}))
.unwrap(); .unwrap();
@@ -586,29 +499,6 @@ mod tests {
"writes do not invent execution passes" "writes do not invent execution passes"
); );
assert!(session.snapshot().unwrap().active.is_empty()); assert!(session.snapshot().unwrap().active.is_empty());
assert_eq!(
session.analytics_output_revision(),
None,
"text fixture files are not classified as成果"
);
let outputs = apply(&root, &json!({"patch":"*** Begin Patch\n*** Add File: game/result.js\n+const result = 1;\n*** Add File: game/style.css\n+body { color: red; }\n*** End Patch"})).await.unwrap();
assert_eq!(outputs["status"], "completed", "{outputs}");
let revision = outputs["revision"].as_u64().unwrap().to_string();
assert_eq!(session.analytics_output_revision(), Some(revision.clone()));
let partial_output = apply(&root, &json!({"patch":"*** Begin Patch\n*** Add File: game/partial.js\n+const partial = 1;\n*** Update File: game/absent.js\n@@\n-old\n+new\n*** End Patch"})).await.unwrap();
assert_eq!(partial_output["status"], "failed");
assert_eq!(session.analytics_output_revision(), Some(revision.clone()));
let events = super::super::direct_tool_bridge::drain_analytics_test_writer(
&config, &context, &writer,
);
let revisions: Vec<_> = events
.iter()
.filter(|event| event["event_name"] == "project_revision_created")
.collect();
assert_eq!(revisions.len(), 1);
assert_eq!(revisions[0]["user_id"], "A");
assert_eq!(revisions[0]["properties"]["revision_id"], revision);
assert_eq!(revisions[0]["properties"]["files_changed_count"], 2);
session.interrupt("fixture stopped".into()).unwrap(); session.interrupt("fixture stopped".into()).unwrap();
assert!(apply( assert!(apply(
&root, &root,
@@ -4243,44 +4243,7 @@ pub(crate) async fn run_direct_browser_evidence_with_cancellation_at(
advisory_interaction: bool, advisory_interaction: bool,
cancellation: Option<Arc<std::sync::atomic::AtomicBool>>, cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
) -> Result<BrowserValidationResult, String> { ) -> Result<BrowserValidationResult, String> {
run_direct_browser_evidence_with_analytics_at(
root,
evidence_root,
scenario,
advisory_interaction,
cancellation,
None,
)
.await
}
pub(crate) async fn run_direct_browser_evidence_with_analytics_at(
root: &Path,
evidence_root: PathBuf,
scenario: Option<crate::browser::BrowserPlaytestScenario>,
advisory_interaction: bool,
cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
) -> Result<BrowserValidationResult, String> {
let observation = crate::analytics::preview::prepare(
root,
capture,
crate::analytics::contract::Source::Direct,
crate::analytics::contract::PreviewSource::Agent,
);
let (_analytics_lease, observation) = match observation {
Some((lease, observation)) => (Some(lease), Some(observation)),
None => (None, None),
};
let (preview, stop_sender) = start_local_game_preview_for_project(root)?; let (preview, stop_sender) = start_local_game_preview_for_project(root)?;
if let Some(observation) = observation {
observation
.with_cancellation(cancellation.clone())
.schedule(preview.port);
}
let validation = crate::browser::validate_local_preview_in_browser_with_cancellation( let validation = crate::browser::validate_local_preview_in_browser_with_cancellation(
BrowserValidationInput { BrowserValidationInput {
url: preview.url, url: preview.url,
@@ -4298,7 +4261,6 @@ pub(crate) async fn run_direct_browser_evidence_with_analytics_at(
cancellation, cancellation,
) )
.await; .await;
drop(_analytics_lease);
let _ = stop_sender.send(()); let _ = stop_sender.send(());
validation validation
} }
@@ -4819,8 +4781,6 @@ pub(crate) async fn run_direct_game_creator_turn_at_with_creation_type(
None, None,
None, None,
None, None,
None,
None,
) )
.await .await
} }
@@ -4832,11 +4792,6 @@ async fn run_direct_game_creator_turn_at_with_creation_type_and_emitter(
turn_emitter: Option<&DirectGameCreatorTurnUpdateEmitter>, turn_emitter: Option<&DirectGameCreatorTurnUpdateEmitter>,
audit: Option<&mut DirectCodexTurnAudit>, audit: Option<&mut DirectCodexTurnAudit>,
direct_user_item: Option<serde_json::Value>, direct_user_item: Option<serde_json::Value>,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
analytics_attempt_id: Option<&str>,
) -> Result<String, String> { ) -> Result<String, String> {
if !root.is_absolute() || !root.is_dir() { if !root.is_absolute() || !root.is_dir() {
return Err("当前项目目录不存在或不是绝对路径".to_string()); return Err("当前项目目录不存在或不是绝对路径".to_string());
@@ -4859,8 +4814,6 @@ async fn run_direct_game_creator_turn_at_with_creation_type_and_emitter(
turn_emitter, turn_emitter,
audit, audit,
direct_user_item, direct_user_item,
capture,
analytics_attempt_id,
) )
.await .await
{ {
@@ -5061,11 +5014,6 @@ async fn run_direct_game_creator_turn_inner(
turn_emitter: Option<&DirectGameCreatorTurnUpdateEmitter>, turn_emitter: Option<&DirectGameCreatorTurnUpdateEmitter>,
audit: Option<&mut DirectCodexTurnAudit>, audit: Option<&mut DirectCodexTurnAudit>,
direct_user_item: Option<serde_json::Value>, direct_user_item: Option<serde_json::Value>,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
analytics_attempt_id: Option<&str>,
) -> Result<String, DirectCodexTurnFailure> { ) -> Result<String, DirectCodexTurnFailure> {
let requires_contract = super::direct_delivery::requires_new_web_contract( let requires_contract = super::direct_delivery::requires_new_web_contract(
root, root,
@@ -5087,49 +5035,16 @@ async fn run_direct_game_creator_turn_inner(
DirectCodexTurnFailure::new(DirectCodexFailureStage::CodeGeneration, error) DirectCodexTurnFailure::new(DirectCodexFailureStage::CodeGeneration, error)
})? })?
.validation; .validation;
let analytics_run = capture.as_ref().map(|(context, _)| { let execution_guard =
crate::analytics::run::Metadata::new( super::direct_execution::begin(root, prompt, requires_contract, execution_config)
context.clone(),
crate::analytics::contract::Source::Direct,
crate::analytics::contract::RunSource::UserSubmit,
)
});
let execution_guard = super::direct_execution::begin(
root,
prompt,
requires_contract,
execution_config,
analytics_run,
)
.await .await
.map_err(|error| DirectCodexTurnFailure::new(DirectCodexFailureStage::CodeGeneration, error))?; .map_err(|error| {
DirectCodexTurnFailure::new(DirectCodexFailureStage::CodeGeneration, error)
})?;
let execution_session = execution_guard.session(); let execution_session = execution_guard.session();
execution_session.set_analytics_capture(capture.clone());
let started = execution_session
.newly_accepted
.then(std::time::Instant::now);
if execution_session.newly_accepted {
if let Ok(ledger) = execution_session.snapshot() {
if let (Some((_, writer)), Some(metadata)) = (&capture, &ledger.analytics_run) {
crate::analytics::run::accepted(writer, root, &ledger.project_id, metadata);
}
}
}
// 在 guard 仍存活时冻结整体结果,避免 Drop 的中断收尾覆盖真实失败原因。
let result: Result<String, DirectCodexTurnFailure> = async {
if let Some(report) = super::direct_delivery::terminal_report(&execution_session) { if let Some(report) = super::direct_delivery::terminal_report(&execution_session) {
return Ok(report); return Ok(report);
} }
if execution_session.newly_accepted {
if let Ok(ledger) = execution_session.snapshot() {
crate::analytics::goal::accepted(
capture.clone(),
root,
&ledger.project_id,
crate::analytics::contract::Source::Direct,
);
}
}
emit_direct_game_creator_progress(root, "codex.turn", "陶泥儿正在处理这条消息"); emit_direct_game_creator_progress(root, "codex.turn", "陶泥儿正在处理这条消息");
if let Some(emitter) = turn_emitter { if let Some(emitter) = turn_emitter {
emitter.emit("running", Some("preparing"), None, None); emitter.emit("running", Some("preparing"), None, None);
@@ -5419,118 +5334,6 @@ async fn run_direct_game_creator_turn_inner(
})?; })?;
} }
Ok(visible_reply) Ok(visible_reply)
}.await;
if let Ok(ledger) = execution_session.snapshot() {
if let (Some(metadata), Some((end_reason, error_code))) = (
&ledger.analytics_run,
direct_analytics_outcome(
ledger.phase,
ledger.requires_contract || ledger.contract.is_some(),
result.is_err(),
execution_session.was_aborted(),
),
) {
let output_revision = execution_session.analytics_output_revision();
crate::analytics::run::direct_finished(
capture,
root,
&ledger.project_id,
metadata,
analytics_attempt_id,
crate::analytics::run::Outcome {
turn_id: Some(ledger.client_turn_id.clone()),
end_reason,
error_code,
duration_ms: started
.and_then(|start| u64::try_from(start.elapsed().as_millis()).ok()),
output_change_detected: output_revision.as_ref().map(|_| true),
revision_id: output_revision,
},
);
}
}
result
}
fn direct_analytics_outcome(
phase: super::direct_execution::ExecutionPhase,
has_contract: bool,
failed: bool,
aborted: bool,
) -> Option<(
crate::analytics::contract::RunEndReason,
Option<crate::analytics::contract::ErrorCode>,
)> {
use super::direct_execution::ExecutionPhase;
use crate::analytics::contract::{ErrorCode, RunEndReason};
if phase == ExecutionPhase::Interrupted || aborted {
return None;
}
if phase == ExecutionPhase::Exhausted {
return Some((RunEndReason::Failed, Some(ErrorCode::RuntimeFailed)));
}
if failed {
// Direct 当前只保留 stage 和展示错误字符串,不从正文猜 Provider 错误类别。
return Some((
RunEndReason::Failed,
Some(ErrorCode::RuntimeErrorUnclassified),
));
}
if phase == ExecutionPhase::Completed || !has_contract {
return Some((RunEndReason::Finished, None));
}
None
}
#[cfg(test)]
mod direct_analytics_tests {
use super::*;
use crate::agent::direct_execution::ExecutionPhase;
use crate::analytics::contract::{ErrorCode, RunEndReason};
#[test]
fn terminal_reports_do_not_turn_exhaustion_or_cancellation_into_success() {
assert_eq!(
direct_analytics_outcome(ExecutionPhase::Exhausted, true, false, false),
Some((RunEndReason::Failed, Some(ErrorCode::RuntimeFailed)))
);
for failed in [true, false] {
assert_eq!(
direct_analytics_outcome(ExecutionPhase::Interrupted, true, failed, false),
None
);
assert_eq!(
direct_analytics_outcome(ExecutionPhase::Working, false, failed, true),
None
);
}
}
#[test]
fn complete_delivery_does_not_hide_later_projection_failure() {
assert_eq!(
direct_analytics_outcome(ExecutionPhase::Completed, true, true, false),
Some((
RunEndReason::Failed,
Some(ErrorCode::RuntimeErrorUnclassified)
))
);
assert_eq!(
direct_analytics_outcome(ExecutionPhase::Completed, true, false, false),
Some((RunEndReason::Finished, None))
);
assert_eq!(
direct_analytics_outcome(ExecutionPhase::Working, false, false, false),
Some((RunEndReason::Finished, None))
);
for phase in [
ExecutionPhase::Working,
ExecutionPhase::Draining,
ExecutionPhase::Sealing,
] {
assert_eq!(direct_analytics_outcome(phase, true, false, false), None);
}
}
} }
#[cfg(test)] #[cfg(test)]
@@ -33,9 +33,7 @@ pub(crate) async fn chat_with_game_creator_direct_codex(
user_item: DirectCodexUserItem, user_item: DirectCodexUserItem,
creation_type: Option<String>, creation_type: Option<String>,
client_turn_id: Option<String>, client_turn_id: Option<String>,
analytics_attempt_id: Option<String>,
) -> Result<String, String> { ) -> Result<String, String> {
let capture = crate::analytics::gui::capture_writer_context();
let root = Path::new(project_path.trim()); let root = Path::new(project_path.trim());
let turn_id = normalize_direct_client_turn_id(client_turn_id.as_deref())?; let turn_id = normalize_direct_client_turn_id(client_turn_id.as_deref())?;
let _active_invocation = DirectTaonierActiveInvocationGuard::enter(root, &turn_id)?; let _active_invocation = DirectTaonierActiveInvocationGuard::enter(root, &turn_id)?;
@@ -64,8 +62,6 @@ pub(crate) async fn chat_with_game_creator_direct_codex(
// DirectProject 的完整回合权威已经落在 project.jsonl;不再创建平行审计日志。 // DirectProject 的完整回合权威已经落在 project.jsonl;不再创建平行审计日志。
None, None,
canonical_user_item, canonical_user_item,
capture,
analytics_attempt_id.as_deref(),
) )
.await .await
{ {
@@ -603,6 +603,30 @@ mod tests {
)); ));
} }
/// 失败终态与正常终态同权:`turn.completed(status="failed")` 必须顶替更早的 `turn.started`
/// 成为锚点,否则队列被回收后新订阅只会看到 `turn.started`,把这轮已收口的回合重放成"还在跑"。
#[test]
fn failed_turn_completed_replaces_started_anchor() {
let mut manager = DirectThreadManager::with_limits(100, 100_000);
manager.append("thread-1", DirectThreadEvent::turn_started(1_000));
manager.append(
"thread-1",
DirectThreadEvent::turn_completed_failed(
crate::agent::DirectTurnFailure::new("host-dropped", "回合宿主任务提前结束"),
FIXED_AT_MS,
),
);
let bootstrap = manager.subscribe("thread-1");
assert!(matches!(
bootstrap.events.as_slice(),
[DirectThreadEvent::TurnCompleted { status, failure, at, .. }]
if status == "failed"
&& failure.as_ref().is_some_and(|failure| failure.kind == "host-dropped")
&& *at == Some(FIXED_AT_MS)
));
}
/// 阶段时间必须随事件一起进队列:bootstrap 与重复订阅都拿到**原值**, /// 阶段时间必须随事件一起进队列:bootstrap 与重复订阅都拿到**原值**,
/// 重放不得重新取钟(否则每次重连都会把已固定的起止时间改掉)。 /// 重放不得重新取钟(否则每次重连都会把已固定的起止时间改掉)。
#[test] #[test]
@@ -211,6 +211,30 @@ impl DirectThreadRequestKind {
} }
} }
/// 失败终态的可下发载荷(`turn.completed.status == "failed"` 时必有,其余终态没有)。
///
/// `kind` 是稳定分类,只给界面选语气,不参与流程分支;`message` 是**已在宿主侧脱敏并截断**的
/// 可展示原因——失败原因只走这一条通道,前端不再从命令返回或另一条 IPC 里另造文案。
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) struct DirectTurnFailure {
/// 稳定失败分类:`timeout` / `model-failed` / `transport-failed` / `request-rejected` /
/// `host-dropped`。
pub(crate) kind: String,
/// 脱敏 + 截断后的失败原因。
pub(crate) message: String,
}
impl DirectTurnFailure {
pub(crate) fn new(kind: impl Into<String>, message: impl Into<String>) -> Self {
Self {
kind: kind.into(),
message: message.into(),
}
}
}
/// Thread Manager 下发的运行态事件。 /// Thread Manager 下发的运行态事件。
/// ///
/// 顺序由数组顺序给出(同一个 subscriber 的 `consume` 按队列顺序返回),因此不需要 `seq`: /// 顺序由数组顺序给出(同一个 subscriber 的 `consume` 按队列顺序返回),因此不需要 `seq`:
@@ -250,7 +274,13 @@ pub(crate) enum DirectThreadEvent {
}, },
#[serde(rename = "turn.completed")] #[serde(rename = "turn.completed")]
TurnCompleted { TurnCompleted {
/// 终态语义:`completed` / `interrupted` / `aborted` 是正常收场;`failed` 是**失败**
/// 此时必须带 `failure` 载荷。
status: String, status: String,
/// 失败载荷:只有 `status == "failed"` 才有;失败原因只从这里下发一次。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional)]
failure: Option<DirectTurnFailure>,
/// 本轮终态的阶段时间(毫秒):宿主处理终态的毫秒钟,或 `durationMs` + 高精度起点的派生值。 /// 本轮终态的阶段时间(毫秒):宿主处理终态的毫秒钟,或 `durationMs` + 高精度起点的派生值。
#[serde(default, skip_serializing_if = "Option::is_none")] #[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<f64>")] #[ts(optional, as = "Option<f64>")]
@@ -301,11 +331,30 @@ impl DirectThreadEvent {
pub(crate) fn turn_completed(status: String, at: u64) -> Self { pub(crate) fn turn_completed(status: String, at: u64) -> Self {
Self::TurnCompleted { Self::TurnCompleted {
status, status,
failure: None,
at: Some(at), at: Some(at),
user_item_id: None, user_item_id: None,
} }
} }
/// 失败终态:`status` 固定 `"failed"`,原因必须随事件一起带出去。
pub(crate) fn turn_completed_failed(failure: DirectTurnFailure, at: u64) -> Self {
Self::TurnCompleted {
status: "failed".to_string(),
failure: Some(failure),
at: Some(at),
user_item_id: None,
}
}
/// 失败载荷:只有失败终态有。
pub(crate) fn failure(&self) -> Option<&DirectTurnFailure> {
match self {
Self::TurnCompleted { failure, .. } => failure.as_ref(),
_ => None,
}
}
/// 附上本轮开口用户条目的 canonical itemId。 /// 附上本轮开口用户条目的 canonical itemId。
/// ///
/// 只在构造之后补一次身份,避免 `turn.started` / `turn.completed` 的既有调用点(含各处兜底 /// 只在构造之后补一次身份,避免 `turn.started` / `turn.completed` 的既有调用点(含各处兜底
@@ -317,8 +366,14 @@ impl DirectThreadEvent {
.map(str::to_string); .map(str::to_string);
match self { match self {
Self::TurnStarted { at, .. } => Self::TurnStarted { at, user_item_id }, Self::TurnStarted { at, .. } => Self::TurnStarted { at, user_item_id },
Self::TurnCompleted { status, at, .. } => Self::TurnCompleted { Self::TurnCompleted {
status, status,
failure,
at,
..
} => Self::TurnCompleted {
status,
failure,
at, at,
user_item_id, user_item_id,
}, },
@@ -1351,4 +1406,56 @@ mod tests {
); );
assert_eq!(item_event.user_item_id(), None); assert_eq!(item_event.user_item_id(), None);
} }
/// 失败终态:`status="failed"` 必须带 `failure{kind,message}`,正常终态不带;载荷跟着身份
/// 一起流转,缺载荷的 `failed` 事件仍能反序列化(前端按"没有原因"处理,不猜)。
#[test]
fn turn_completed_carries_failure_payload_only_when_failed() {
let failed = DirectThreadEvent::turn_completed_failed(
DirectTurnFailure::new("model-failed", "上游返回 500:模型服务暂不可用"),
4_000,
)
.with_user_item_id(Some("direct-codex:turn-1:user"));
assert_eq!(
failed.failure(),
Some(&DirectTurnFailure::new(
"model-failed",
"上游返回 500:模型服务暂不可用"
))
);
assert_eq!(failed.user_item_id(), Some("direct-codex:turn-1:user"));
assert_eq!(
serde_json::to_value(&failed).expect("serialize failed turn"),
json!({
"type": "turn.completed",
"status": "failed",
"failure": {"kind": "model-failed", "message": "上游返回 500:模型服务暂不可用"},
"at": 4_000u64,
"userItemId": "direct-codex:turn-1:user",
})
);
assert_eq!(
serde_json::from_value::<DirectThreadEvent>(
serde_json::to_value(&failed).expect("serialize")
)
.expect("round trip"),
failed
);
// 正常终态不带载荷,也不回写 `failure: null`。
let completed = DirectThreadEvent::turn_completed("completed".to_string(), 5_000);
assert_eq!(completed.failure(), None);
assert_eq!(
serde_json::to_value(&completed).expect("serialize completed turn"),
json!({"type": "turn.completed", "status": "completed", "at": 5_000u64})
);
// 精简 / 旧形状:`failed` 但没有载荷也要能反序列化。
let sparse: DirectThreadEvent = serde_json::from_value(json!({
"type": "turn.completed",
"status": "failed",
}))
.expect("failed turn without failure payload");
assert_eq!(sparse.failure(), None);
}
} }
@@ -1600,15 +1600,6 @@ fn bridge_write_file(root: &Path, arguments: &Value) -> Value {
bridge_write_file_with_permit(root, arguments, None) bridge_write_file_with_permit(root, arguments, None)
} }
fn bridge_file_content_changed(root: &Path, path: &str, content: &[u8]) -> Option<bool> {
crate::analytics::project::file_content_changed(
root,
path,
content,
DIRECT_TOOL_BRIDGE_MAX_WRITE_CONTENT_BYTES as u64,
)
}
fn bridge_write_file_with_permit( fn bridge_write_file_with_permit(
root: &Path, root: &Path,
arguments: &Value, arguments: &Value,
@@ -1650,11 +1641,6 @@ fn bridge_write_file_with_permit(
"direct-codex.file.write", "direct-codex.file.write",
)?; )?;
let lock_wait_ms = acquire_started.elapsed().as_millis(); let lock_wait_ms = acquire_started.elapsed().as_millis();
let analytics_change_kind = write_permit.and_then(|_| {
let kind = crate::analytics::project::file_change_kind(&path)?;
(bridge_file_content_changed(root, &path, content.as_bytes()) == Some(true))
.then_some(kind)
});
let write_started = std::time::Instant::now(); let write_started = std::time::Instant::now();
let commit = || { let commit = || {
let written = write_local_project_file_at(root, &path, content)?; let written = write_local_project_file_at(root, &path, content)?;
@@ -1667,9 +1653,6 @@ fn bridge_write_file_with_permit(
Some(permit) => permit.run(commit)?, Some(permit) => permit.run(commit)?,
None => commit()?, None => commit()?,
}; };
if let (Some(permit), Some(kind)) = (write_permit, analytics_change_kind) {
permit.record_analytics_revision(revision, kind, 1);
}
// 现场一次 2.6KB 写入实测 5.5 秒。只在明显偏慢时记账,正常写入不刷日志。 // 现场一次 2.6KB 写入实测 5.5 秒。只在明显偏慢时记账,正常写入不刷日志。
if lock_wait_ms + write_ms > 200 { if lock_wait_ms + write_ms > 200 {
app_log!( app_log!(
@@ -3490,82 +3473,6 @@ pub(in crate::agent) async fn generate_images_concurrently_for_test(
.await .await
} }
#[cfg(test)]
pub(super) fn analytics_test_writer(
config: &Path,
) -> (
crate::analytics::run::Metadata,
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
) {
use crate::analytics::{
contract::{Context, Route, RunSource, Source},
run,
store::AnalyticsWriter,
};
let mut context = Context {
route: Route::from_identity(Some("A".into()), Some("https://example.com")),
editor_session_id: uuid::Uuid::new_v4().to_string(),
client_version: "1.0.0".into(),
};
fs::create_dir_all(config).expect("create analytics test config directory");
let metadata = run::Metadata::new(context.clone(), Source::Direct, RunSource::UserSubmit);
context.route.user_id = Some("B".into());
let writer = AnalyticsWriter::start(config.into(), context.editor_session_id.clone());
(metadata, context, writer)
}
#[cfg(test)]
pub(super) fn drain_analytics_test_writer(
config: &Path,
context: &crate::analytics::contract::Context,
writer: &crate::analytics::store::AnalyticsWriter,
) -> Vec<Value> {
use crate::analytics::contract::{EntrySource, EventData, SessionStart, Source};
let marker = context
.capture(
EventData::EditorSessionStart(SessionStart {
entry_source: EntrySource::DirectLaunch,
first_project_id: None,
}),
None,
Source::Editor,
None,
)
.unwrap();
let marker_id = marker.event_id.clone();
assert!(writer.try_record(context.route.clone(), marker, marker_id.clone()));
assert!(writer.flush());
let batches = config
.join("analytics/instances")
.join(&context.editor_session_id)
.join("batches");
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
let events: Vec<Value> = fs::read_dir(&batches)
.into_iter()
.flatten()
.flatten()
.filter(|entry| !entry.file_name().to_string_lossy().starts_with('.'))
.filter_map(|entry| fs::read_to_string(entry.path().join("events.jsonl")).ok())
.flat_map(|contents| {
contents
.lines()
.filter_map(|line| serde_json::from_str::<Value>(line).ok())
.collect::<Vec<_>>()
})
.collect();
if events.iter().any(|event| event["event_id"] == marker_id) {
return events;
}
assert!(
std::time::Instant::now() < deadline,
"analytics FIFO sentinel timed out"
);
std::thread::sleep(std::time::Duration::from_millis(5));
}
}
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
#[tokio::test] #[tokio::test]
@@ -4355,293 +4262,6 @@ mod tests {
assert_eq!(importability.get("assets/vector.svg"), Some(&true)); assert_eq!(importability.get("assets/vector.svg"), Some(&true));
} }
#[tokio::test]
async fn analytics_real_file_write_preserves_original_identity_and_failed_run_revision() {
use crate::analytics::{
contract::{ErrorCode, RunEndReason},
run,
};
let temporary = tempfile::tempdir().unwrap();
let root = temporary.path().join("project");
let config = temporary.path().join("config");
let (metadata, context, writer) = analytics_test_writer(&config);
let original_capture = Some((metadata.context.clone(), writer.clone()));
let mut lifecycle = crate::analytics::gui::LifecycleFixture::start(
metadata.context.clone(),
writer.clone(),
);
lifecycle.create_and_open(&root, "direct-analytics");
let session = super::super::direct_execution::open_with_analytics_at(
&temporary.path().join("host"),
&root,
"analytics-write",
&format!("{:x}", Sha256::digest(b"request")),
false,
&Default::default(),
Some(metadata.clone()),
)
.unwrap();
session.set_analytics_capture(Some((context.clone(), writer.clone())));
super::super::direct_delivery::register_contract(
&root,
&session,
&json!({
"scope": "核对当前项目宿主写入的成果采集",
"changeKind": "project",
"requirements": [{"id": "analytics-output", "kind": "artifact", "path": "game/index.html"}]
}),
)
.await
.expect("freeze a validated delivery contract before writing");
let project_id = session.snapshot().unwrap().project_id;
run::accepted(&writer, &root, &project_id, &metadata);
crate::analytics::goal::accepted(
original_capture.clone(),
&root,
&project_id,
crate::analytics::contract::Source::Direct,
);
let lease = session
.admit(super::super::direct_execution::EffectKind::Write, None)
.unwrap();
let permit = lease.write_permit().unwrap();
let arguments = json!({"path":"game/index.html", "content":"<!doctype html><html><body>真实预览</body></html>"});
let changed = bridge_write_file_with_permit(&root, &arguments, Some(&permit));
assert_eq!(changed["isError"], false);
let payload: Value =
serde_json::from_str(changed["content"][0]["text"].as_str().unwrap()).unwrap();
let revision = payload["revision"].as_u64().unwrap().to_string();
assert_eq!(session.analytics_output_revision(), Some(revision.clone()));
assert_eq!(
bridge_write_file_with_permit(&root, &arguments, Some(&permit))["isError"],
false
);
assert_eq!(
bridge_write_file_with_permit(
&root,
&json!({"path":"../bad.js","content":"bad"}),
Some(&permit)
)["isError"],
true
);
assert_eq!(session.analytics_output_revision(), Some(revision.clone()));
lease.finish(true, true, None).unwrap();
let attempt = uuid::Uuid::new_v4().to_string();
run::direct_finished(
Some((context.clone(), writer.clone())),
&root,
&project_id,
&metadata,
Some(&attempt),
run::Outcome {
turn_id: Some("analytics-write".into()),
end_reason: RunEndReason::Failed,
error_code: Some(ErrorCode::RuntimeErrorUnclassified),
duration_ms: None,
output_change_detected: Some(true),
revision_id: session.analytics_output_revision(),
},
);
run::settle(Some((context.clone(), writer.clone())), &attempt, false);
// 默认项目使用 npm:预览服务读取真实构建目录,夹具提供构建入口,不调用构建器。
let served_root = crate::project_game_root(&root);
fs::create_dir_all(&served_root).unwrap();
fs::write(
served_root.join("index.html"),
"<!doctype html><html><body>真实预览构建</body></html>",
)
.unwrap();
let (preview_lease, observation) = crate::analytics::preview::prepare(
&root,
original_capture.clone(),
crate::analytics::contract::Source::Editor,
crate::analytics::contract::PreviewSource::User,
)
.unwrap();
let (preview, stop) = crate::start_local_game_preview_for_project(&root).unwrap();
observation.observe(preview.port).await;
let preview_revision = crate::read_game_creator_agent_runtime_project_revision(&root)
.unwrap()
.revision
.to_string();
drop(preview_lease);
let _ = stop.send(());
let checkpoint = crate::commands::checkpoint_with_capture_for_test(
root.to_string_lossy().into_owned(),
original_capture,
)
.unwrap();
lifecycle.exit();
let events = drain_analytics_test_writer(&config, &context, &writer);
let ids: std::collections::HashSet<_> = events
.iter()
.map(|event| event["event_id"].as_str().unwrap())
.collect();
assert_eq!(ids.len(), events.len());
let chain: Vec<_> = events
.iter()
.filter(|event| event["user_id"] == "A")
.collect();
for name in [
"editor_session_start",
"editor_focus_start",
"project_create_success",
"project_open",
"creative_task_submit",
"project_revision_created",
"agent_run_failed",
"preview_ready",
"project_save",
"editor_focus_end",
"editor_session_end",
] {
assert_eq!(
chain
.iter()
.filter(|event| event["event_name"] == name)
.count(),
1,
"{name}"
);
}
for event in &chain {
assert_eq!(
event["editor_session_id"],
metadata.context.editor_session_id
);
if !event["project_id"].is_null() {
assert_eq!(event["project_id"], project_id);
}
if !event["creative_task_id"].is_null() {
assert_eq!(event["creative_task_id"], project_id);
}
if !event["agent_run_id"].is_null() {
assert_eq!(event["agent_run_id"], metadata.run_id);
}
let name = event["event_name"].as_str().unwrap();
if matches!(
name,
"project_create_success"
| "project_open"
| "creative_task_submit"
| "project_revision_created"
| "agent_run_failed"
| "preview_ready"
| "project_save"
) {
assert_eq!(
event["project_id"], project_id,
"{name} must identify its project"
);
}
if matches!(
name,
"creative_task_submit"
| "project_revision_created"
| "agent_run_failed"
| "preview_ready"
| "project_save"
) {
assert_eq!(
event["creative_task_id"], project_id,
"{name} must identify its goal"
);
}
if name == "agent_run_failed" {
assert_eq!(event["agent_run_id"], metadata.run_id);
assert_eq!(event["agent_turn_id"], "analytics-write");
}
}
let save = chain
.iter()
.find(|event| event["event_name"] == "project_save")
.unwrap();
assert_eq!(save["properties"]["save_source"], "checkpoint");
assert!(Path::new(&checkpoint.checkpoint_path).is_dir());
let mut digest = Sha256::new();
digest.update(serde_json::to_vec(&metadata.context.route).unwrap());
digest.update([0]);
digest.update(format!("{project_id}:{}:project_save", checkpoint.checkpoint_id).as_bytes());
let checkpoint_fact = format!("{:x}", digest.finalize());
let batches = config
.join("analytics/instances")
.join(&context.editor_session_id)
.join("batches");
let checkpoint_fact_matches = fs::read_dir(batches)
.unwrap()
.flatten()
.filter_map(|entry| fs::read(entry.path().join("meta.json")).ok())
.filter_map(|bytes| serde_json::from_slice::<Value>(&bytes).ok())
.filter(|batch| batch["facts"][&checkpoint_fact] == save["event_id"])
.count();
assert_eq!(
checkpoint_fact_matches, 1,
"save fact must use the actual checkpoint ID"
);
let ready = chain
.iter()
.find(|event| event["event_name"] == "preview_ready")
.unwrap();
assert_eq!(ready["properties"]["preview_version"], preview_revision);
let focus_start = chain
.iter()
.find(|event| event["event_name"] == "editor_focus_start")
.unwrap();
let focus_end = chain
.iter()
.find(|event| event["event_name"] == "editor_focus_end")
.unwrap();
assert_eq!(
focus_start["properties"]["focus_interval_id"],
focus_end["properties"]["focus_interval_id"]
);
let revisions: Vec<_> = events
.iter()
.filter(|event| event["event_name"] == "project_revision_created")
.collect();
assert_eq!(revisions.len(), 1);
assert_eq!(revisions[0]["user_id"], "A");
assert_eq!(revisions[0]["properties"]["revision_id"], revision);
let failed: Vec<_> = events
.iter()
.filter(|event| event["event_name"] == "agent_run_failed")
.collect();
assert_eq!(failed.len(), 1);
assert_eq!(failed[0]["user_id"], "A");
assert_eq!(failed[0]["properties"]["revision_id"], revision);
assert_eq!(failed[0]["properties"]["output_change_detected"], true);
}
#[test]
fn analytics_file_comparison_requires_known_bounded_content() {
let temporary = tempfile::tempdir().unwrap();
let root = temporary.path();
assert_eq!(
bridge_file_content_changed(root, "new.js", b"new"),
Some(true)
);
fs::write(root.join("new.js"), b"new").unwrap();
assert_eq!(
bridge_file_content_changed(root, "new.js", b"new"),
Some(false)
);
assert_eq!(
bridge_file_content_changed(root, "new.js", b"changed"),
Some(true)
);
fs::create_dir(root.join("directory.js")).unwrap();
assert_eq!(
bridge_file_content_changed(root, "directory.js", b"new"),
None
);
fs::write(
root.join("large.js"),
vec![0; DIRECT_TOOL_BRIDGE_MAX_WRITE_CONTENT_BYTES + 1],
)
.unwrap();
assert_eq!(bridge_file_content_changed(root, "large.js", b"new"), None);
}
#[test] #[test]
fn bridge_write_file_writes_project_relative_text_without_runtime_tasks() { fn bridge_write_file_writes_project_relative_text_without_runtime_tasks() {
let temporary = tempfile::tempdir().expect("create direct write root"); let temporary = tempfile::tempdir().expect("create direct write root");

Some files were not shown because too many files have changed in this diff Show More