Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,10 @@ jobs:
if: github.event_name == 'pull_request'
needs: [test]
runs-on: ubuntu-latest
permissions:
contents: read
issues: write
pull-requests: write
steps:
- uses: actions/checkout@v4
- uses: actions/download-artifact@v4
Expand Down
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ apps/web/dist/

# Dependencies
node_modules/
.data/

# Env
.env
Expand Down
128 changes: 103 additions & 25 deletions apps/api/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,15 +3,46 @@ import { streamSSE } from "hono/streaming";
import type { BattleEvent } from "@agent-arena/contracts";
import fixture from "../../../examples/fixtures/hackathon-001.json";
import verifiedShowcase from "../../../examples/fixtures/verified-showcase.json";
import { runLiveBattleFromPayload, LiveBattleIdeaTooLongError } from "@/lib/runtime/runLiveBattleFromPayload";
import { StepFunNotConfiguredError } from "@/lib/runtime/providers/stepfun";
import { runLiveBattleFromPayload } from "@/lib/runtime/runLiveBattleFromPayload";
import { BattleRateLimiter } from "./middlewares/rate-limit";
import { isLiveBattleEnabled } from "./middlewares/feature-flag";
import { LocalLiveBattleStore } from "./live-battle-store";

export const app = new Hono();

const liveBattleRateLimiter = new BattleRateLimiter();
const liveBattlesInFlight = new Map<string, { idea: string; startedAt: number }>();
const liveBattleStore = new LocalLiveBattleStore();
type LiveMessage = { type: "battle"; event: BattleEvent } | { type: "done" } | { type: "error"; error: string };
const liveSubscribers = new Map<string, Set<(message: LiveMessage) => Promise<void>>>();
const liveRuns = new Map<string, Promise<void>>();

async function publishLiveMessage(battleId: string, message: LiveMessage): Promise<void> {
const subscribers = [...(liveSubscribers.get(battleId) ?? [])];
await Promise.allSettled(subscribers.map((subscriber) => subscriber(message)));
}

function startLiveBattle(battleId: string, idea: string): Promise<void> {
const existing = liveRuns.get(battleId);
if (existing) return existing;
const run = (async () => {
try {
for await (const event of runLiveBattleFromPayload({ battleId, idea })) {
await liveBattleStore.append(battleId, event);
await publishLiveMessage(battleId, { type: "battle", event });
}
await liveBattleStore.finish(battleId, "completed");
await publishLiveMessage(battleId, { type: "done" });
} catch (error) {
const message = error instanceof Error ? error.message : "实时战斗异常中断";
await liveBattleStore.finish(battleId, "failed", message);
await publishLiveMessage(battleId, { type: "error", error: message });
} finally {
liveRuns.delete(battleId);
}
})();
liveRuns.set(battleId, run);
return run;
}

app.get("/api/health", (context) =>
context.json({ status: "ok", service: "agent-arena-api" }),
Expand Down Expand Up @@ -113,6 +144,18 @@ app.get("/api/battles/:id/events", async (context) => {
return context.json({ battleId, source: "fixture", events: demoEvents() });
}

const localBattle = await liveBattleStore.get(battleId);
if (localBattle) {
return context.json({ battleId, source: "local-event-store", status: localBattle.status, error: localBattle.error, events: localBattle.events });
}

// Avoid loading the full Postgres adapter when persistence is not configured.
// The documented fallback is immediate and must not spend the request budget
// compiling database drivers only to discover DATABASE_URL is absent.
if (!process.env.DATABASE_URL) {
return context.json({ battleId, source: "fallback", events: [] });
}

// Postgres/event-store integration is deliberately soft-failing: the replay
// remains usable and never blocks the battle experience when storage is absent.
try {
Expand Down Expand Up @@ -179,10 +222,10 @@ app.post("/api/battles", async (context) => {
return context.json({ error: "创意不能为空" }, 400);
}

const ip =
context.req.header("cf-connecting-ip") ??
context.req.header("x-forwarded-for")?.split(",")[0]?.trim() ??
"anonymous";
const trustProxyHeaders = process.env.TRUST_PROXY_HEADERS === "true";
const ip = trustProxyHeaders
? context.req.header("cf-connecting-ip") ?? context.req.header("x-forwarded-for")?.split(",")[0]?.trim() ?? "proxied-client"
: "direct-client";
const decision = liveBattleRateLimiter.check(ip);
if (!decision.allowed) {
return context.json(
Expand All @@ -193,7 +236,7 @@ app.post("/api/battles", async (context) => {
}

const battleId = `live_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 8)}`;
liveBattlesInFlight.set(battleId, { idea, startedAt: Date.now() });
await liveBattleStore.create(battleId, idea);
return context.json({ battleId, sseUrl: `/api/battles/${battleId}/stream` }, 201);
});

Expand All @@ -202,38 +245,73 @@ app.get("/api/battles/:id/stream", async (context) => {
return context.json({ error: "实时 AI 竞技当前未开启" }, 501);
}
const battleId = context.req.param("id");
const pending = liveBattlesInFlight.get(battleId);
if (!pending) {
const stored = await liveBattleStore.get(battleId);
if (!stored) {
return context.json({ error: "找不到该战斗" }, 404);
}

return streamSSE(context, async (stream) => {
const heartbeatMs = Number.parseInt(process.env.SSE_HEARTBEAT_MS ?? "15000", 10);
const configuredHeartbeatMs = Number.parseInt(process.env.SSE_HEARTBEAT_MS ?? "2000", 10);
const heartbeatMs = Math.min(3000, Math.max(1000, Number.isFinite(configuredHeartbeatMs) ? configuredHeartbeatMs : 2000));
let closed = false;
const heartbeat = setInterval(() => {
if (closed) return;
void stream.writeSSE({ event: "heartbeat", data: JSON.stringify({ at: Date.now() }) });
}, heartbeatMs);

const sentIds = new Set<string>();
let writeQueue = Promise.resolve();
const write = (event: string, data: unknown): Promise<void> => {
writeQueue = writeQueue.then(() => stream.writeSSE({ event, data: JSON.stringify(data) }));
return writeQueue;
};
let resolveStream: () => void = () => undefined;
const completed = new Promise<void>((resolve) => { resolveStream = resolve; });
const subscriber = async (message: LiveMessage): Promise<void> => {
if (message.type === "battle") {
if (sentIds.has(message.event.id)) return;
sentIds.add(message.event.id);
await write("battle", message.event);
return;
}
if (message.type === "done") await write("done", { battleId });
else await write("error", { error: message.error });
liveSubscribers.get(battleId)?.delete(subscriber);
resolveStream();
};

try {
for await (const event of runLiveBattleFromPayload({ battleId, idea: pending.idea })) {
await stream.writeSSE({ event: "battle", data: JSON.stringify(event) });
const subscribers = liveSubscribers.get(battleId) ?? new Set();
subscribers.add(subscriber);
liveSubscribers.set(battleId, subscribers);

const latest = await liveBattleStore.get(battleId);
if (!latest) return;
for (const event of latest.events) {
if (sentIds.has(event.id)) continue;
sentIds.add(event.id);
await write("battle", event);
}
if (latest.status === "completed") {
await write("done", { battleId });
return;
}
if (latest.status === "failed") {
await write("error", { error: latest.error ?? "实时战斗异常中断" });
return;
}
if (latest.status === "running" && !liveRuns.has(battleId)) {
const interruption = "服务恢复后检测到未完成战斗;已保留现有证据,请重新发起以避免重复运行智能体。";
await liveBattleStore.finish(battleId, "failed", interruption);
await write("error", { error: interruption });
return;
}
await stream.writeSSE({ event: "done", data: JSON.stringify({ battleId }) });
} catch (err) {
const message =
err instanceof StepFunNotConfiguredError
? err.message
: err instanceof LiveBattleIdeaTooLongError
? err.message
: err instanceof Error
? err.message
: "实时战斗异常中断";
await stream.writeSSE({ event: "error", data: JSON.stringify({ error: message }) });
void startLiveBattle(battleId, latest.idea);
await completed;
} finally {
liveSubscribers.get(battleId)?.delete(subscriber);
closed = true;
clearInterval(heartbeat);
liveBattlesInFlight.delete(battleId);
}
});
});
21 changes: 20 additions & 1 deletion apps/api/src/index.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,24 @@
import { serve } from "@hono/node-server";
import { app } from "./app";
import { existsSync } from "node:fs";
import path from "node:path";

const envCandidates: string[] = [];
let envDirectory = path.resolve(".");
while (true) {
envCandidates.push(path.join(envDirectory, ".env.local"));
const parent = path.dirname(envDirectory);
if (parent === envDirectory) break;
envDirectory = parent;
}

for (const candidate of envCandidates) {
if (existsSync(candidate)) {
process.loadEnvFile(candidate);
break;
}
}

const { app } = await import("./app");

serve({ fetch: app.fetch, port: 8787 }, (info) => {
console.log(`Agent Arena API listening on http://localhost:${info.port}`);
Expand Down
41 changes: 41 additions & 0 deletions apps/api/src/live-battle-store.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
import { mkdtemp, readFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import path from "node:path";
import { describe, expect, it } from "vitest";
import type { BattleEvent } from "@agent-arena/contracts";
import { LocalLiveBattleStore, resolveLiveBattleRoot } from "./live-battle-store";

const event: BattleEvent = {
id: "event_1", battleId: "live_test", round: "briefing", eventType: "brief_created",
title: "简报下发", content: "测试", createdAt: "2026-07-25T00:00:00.000Z",
};

describe("LocalLiveBattleStore", () => {
it("resolves the default data directory from the API module, not process.cwd", () => {
const root = resolveLiveBattleRoot({} as NodeJS.ProcessEnv, "file:///C:/workspace/apps/api/src/live-battle-store.ts");
expect(root.replaceAll("\\", "/")).toBe("C:/workspace/apps/api/.data/live-battles/");
});
it("persists battle events and terminal status across store instances", async () => {
const root = await mkdtemp(path.join(tmpdir(), "agent-arena-store-"));
const writer = new LocalLiveBattleStore(root);
await writer.create("live_test", "测试创意");
await writer.append("live_test", event);
await writer.finish("live_test", "completed");

const reader = new LocalLiveBattleStore(root);
const restored = await reader.get("live_test");
expect(restored).toMatchObject({ battleId: "live_test", idea: "测试创意", status: "completed" });
expect(restored?.events).toHaveLength(1);
expect(restored?.events[0]).toMatchObject({ id: "event_1", sequence: 1 });
expect(JSON.parse(await readFile(path.join(root, "live_test.json"), "utf8"))).toMatchObject({ status: "completed" });
});

it("deduplicates repeated event ids", async () => {
const root = await mkdtemp(path.join(tmpdir(), "agent-arena-store-"));
const store = new LocalLiveBattleStore(root);
await store.create("live_test", "测试创意");
await store.append("live_test", event);
await store.append("live_test", event);
expect((await store.get("live_test"))?.events).toHaveLength(1);
});
});
96 changes: 96 additions & 0 deletions apps/api/src/live-battle-store.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
import { mkdir, readFile, rename, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import path from "node:path";
import { fileURLToPath } from "node:url";
import type { BattleEvent } from "@agent-arena/contracts";
import { assertBattleEvent } from "@/arena/schemas/validators";

export type LiveBattleStatus = "created" | "running" | "completed" | "failed";

export type StoredLiveBattle = {
battleId: string;
idea: string;
status: LiveBattleStatus;
events: BattleEvent[];
createdAt: string;
updatedAt: string;
error?: string;
};

export function resolveLiveBattleRoot(env: NodeJS.ProcessEnv = process.env, moduleUrl = import.meta.url): string {
if (env.VITEST) return path.join(tmpdir(), `agent-arena-live-${process.pid}`);
if (env.AGENT_ARENA_DATA_DIR) return path.resolve(env.AGENT_ARENA_DATA_DIR);
return fileURLToPath(new URL("../.data/live-battles/", moduleUrl));
}

export class LocalLiveBattleStore {
private readonly root: string;
private readonly queues = new Map<string, Promise<void>>();

constructor(root = resolveLiveBattleRoot()) { this.root = root; }

async create(battleId: string, idea: string): Promise<StoredLiveBattle> {
const at = new Date().toISOString();
const battle: StoredLiveBattle = { battleId, idea, status: "created", events: [], createdAt: at, updatedAt: at };
await this.write(battle);
return battle;
}

async get(battleId: string): Promise<StoredLiveBattle | null> {
try {
return JSON.parse(await readFile(this.fileFor(battleId), "utf8")) as StoredLiveBattle;
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return null;
throw error;
}
}

async append(battleId: string, event: BattleEvent): Promise<void> {
const validatedEvent = {
...event,
actorType: event.actorId ? (event.actorId === "judge_panel" ? "judge" : "agent") : "system",
} as BattleEvent;
assertBattleEvent(validatedEvent);
await this.serial(battleId, async () => {
const battle = await this.get(battleId);
if (!battle) throw new Error(`找不到本地战斗 ${battleId}`);
if (battle.events.some((candidate) => candidate.id === event.id)) return;
battle.events.push({ ...validatedEvent, sequence: battle.events.length + 1 });
battle.status = "running";
battle.updatedAt = new Date().toISOString();
await this.write(battle);
});
}

async finish(battleId: string, status: "completed" | "failed", error?: string): Promise<void> {
await this.serial(battleId, async () => {
const battle = await this.get(battleId);
if (!battle) return;
battle.status = status;
battle.error = error;
battle.updatedAt = new Date().toISOString();
await this.write(battle);
});
}

private fileFor(battleId: string): string {
return path.join(this.root, `${battleId.replace(/[^a-zA-Z0-9_-]/g, "_")}.json`);
}

private async write(battle: StoredLiveBattle): Promise<void> {
await mkdir(this.root, { recursive: true });
const target = this.fileFor(battle.battleId);
const temporary = `${target}.${process.pid}.tmp`;
await writeFile(temporary, JSON.stringify(battle, null, 2), "utf8");
await rename(temporary, target);
}

private async serial(battleId: string, operation: () => Promise<void>): Promise<void> {
const previous = this.queues.get(battleId) ?? Promise.resolve();
const current = previous.then(operation, operation);
this.queues.set(battleId, current);
try { await current; } finally {
if (this.queues.get(battleId) === current) this.queues.delete(battleId);
}
}
}
2 changes: 1 addition & 1 deletion apps/api/src/middlewares/feature-flag.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@

export function isLiveBattleEnabled(env: NodeJS.ProcessEnv = process.env): boolean {
const value = env.AGENT_ARENA_LIVE_BATTLE_ENABLED;
if (typeof value !== "string") return false;
if (typeof value !== "string") return typeof env.STEPFUN_API_KEY === "string" && env.STEPFUN_API_KEY.trim().length > 0;
const normalized = value.trim().toLowerCase();
return normalized === "true" || normalized === "1" || normalized === "yes";
}
Loading
Loading