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
Original file line number Diff line number Diff line change
@@ -1,321 +1,10 @@
import type { FastifyInstance, FastifyReply, FastifyRequest } from "fastify";
import type { Pool } from "pg";

import type { AgentManager } from "../agents/manager.js";
import {
computeActivityStats,
computeDailyStatus,
computeWorkingTimeByProject,
type ActivityEventRow,
} from "../activity-metrics.js";

type ActivityRouteDeps = {
pool: Pool;
agentManager: AgentManager;
parseActivityQuery: (query: Record<string, unknown>) => {
start: Date | null;
end: Date | null;
tz: string;
granularity: "hour" | "day" | "week" | "month";
};
loadScopedActivityEvents: (
aq: ReturnType<ActivityRouteDeps["parseActivityQuery"]>
) => Promise<{ rows: ActivityEventRow[]; rangeStart: Date | null }>;
timeRangeClause: (
aq: ReturnType<ActivityRouteDeps["parseActivityQuery"]>,
column: string,
paramOffset?: number
) => { clause: string; params: unknown[] };
dateTruncTz: (
granularity: "hour" | "day" | "week" | "month",
column: string,
tz: string
) => string;
escapeLike: (s: string) => string;
};

async function handleHeatmap(deps: ActivityRouteDeps, request: FastifyRequest) {
const query = request.query as Record<string, unknown>;
const days = Math.min(
Math.max(parseInt((query.days as string) ?? "365", 10) || 365, 1),
730
);
const aq = deps.parseActivityQuery(query);

const result = await deps.pool.query<{ day: string; count: number }>(
`SELECT ${deps.dateTruncTz("day", "created_at", aq.tz)} AS day, COUNT(*)::int AS count
FROM agent_events
WHERE created_at >= NOW() - make_interval(days => $1)
GROUP BY day ORDER BY day`,
[days]
);

return { days: result.rows };
}

async function handleStats(deps: ActivityRouteDeps, request: FastifyRequest) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const { rows, rangeStart } = await deps.loadScopedActivityEvents(aq);
const eventFilter = deps.timeRangeClause(aq, "created_at");

const busiestDayResult = await deps.pool.query<{
day: string;
count: number;
}>(
`SELECT ${deps.dateTruncTz("day", "created_at", aq.tz)} AS day, COUNT(*)::int AS count
FROM agent_events
${eventFilter.clause}
GROUP BY day ORDER BY count DESC LIMIT 1`,
eventFilter.params
);
const stats = computeActivityStats(rows, rangeStart);

return {
totalWorkingMs: stats.totalWorkingMs,
avgBlockedMs: stats.avgBlockedMs,
avgWaitingMs: stats.avgWaitingMs,
busiestDay: busiestDayResult.rows[0]?.day ?? null,
busiestDayCount: busiestDayResult.rows[0]?.count ?? 0,
stateDurations: stats.stateDurations,
};
}

async function handleDailyStatus(
deps: ActivityRouteDeps,
request: FastifyRequest
) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const { rows, rangeStart } = await deps.loadScopedActivityEvents(aq);
return {
days: computeDailyStatus(rows, rangeStart, aq.granularity),
granularity: aq.granularity,
};
}

async function handleActiveHours(
deps: ActivityRouteDeps,
request: FastifyRequest
) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const eventFilter = deps.timeRangeClause(aq, "created_at");
const result = await deps.pool.query<{ created_at: string }>(
`SELECT created_at::text AS created_at
FROM agent_events
${eventFilter.clause ? `${eventFilter.clause} AND` : "WHERE"} event_type IN ('working', 'blocked', 'waiting_user')
ORDER BY created_at`,
eventFilter.params
);
return { events: result.rows };
}

async function handleAgentsCreated(
deps: ActivityRouteDeps,
request: FastifyRequest
) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const eventFilter = deps.timeRangeClause(aq, "first_seen");
const result = await deps.pool.query<{ day: string; count: number }>(
`SELECT ${deps.dateTruncTz(aq.granularity, "first_seen", aq.tz)} AS day, COUNT(*)::int AS count
FROM (
SELECT agent_id, MIN(created_at) AS first_seen
FROM agent_events
GROUP BY agent_id
) per_agent
${eventFilter.clause}
GROUP BY day ORDER BY day`,
eventFilter.params
);
const total = result.rows.reduce((sum, row) => sum + row.count, 0);
return { days: result.rows, total, granularity: aq.granularity };
}

async function handleWorkingTimeByProject(
deps: ActivityRouteDeps,
request: FastifyRequest
) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const rangeStart = aq.start;
const eventFilter = deps.timeRangeClause(aq, "ae.created_at");
const inRangeResult = await deps.pool.query<ActivityEventRow>(
`SELECT ae.agent_id, ae.event_type, ae.created_at,
COALESCE(ae.project_dir, a.cwd) AS project_dir
FROM agent_events ae
LEFT JOIN agents a ON a.id = ae.agent_id
${eventFilter.clause}
ORDER BY ae.agent_id, ae.created_at`,
eventFilter.params
);

let rows = inRangeResult.rows;
if (rangeStart) {
const boundaryResult = await deps.pool.query<ActivityEventRow>(
`SELECT DISTINCT ON (ae.agent_id) ae.agent_id, ae.event_type, ae.created_at,
COALESCE(ae.project_dir, a.cwd) AS project_dir
FROM agent_events ae
LEFT JOIN agents a ON a.id = ae.agent_id
WHERE ae.created_at < $1
ORDER BY ae.agent_id, ae.created_at DESC`,
[rangeStart]
);
rows = [...boundaryResult.rows, ...inRangeResult.rows].sort((a, b) => {
const agentCompare = a.agent_id.localeCompare(b.agent_id);
if (agentCompare !== 0) return agentCompare;
return a.created_at.getTime() - b.created_at.getTime();
});
}

return { projects: computeWorkingTimeByProject(rows, rangeStart) };
}

async function handleTokenStats(
deps: ActivityRouteDeps,
request: FastifyRequest
) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const tokenFilter = deps.timeRangeClause(
aq,
"COALESCE(session_start, harvested_at)"
);
const result = await deps.pool.query<{
total_input: number;
total_cache_creation: number;
total_cache_read: number;
total_output: number;
total_messages: number;
total_sessions: number;
}>(
`SELECT
COALESCE(SUM(input_tokens), 0) AS total_input,
COALESCE(SUM(cache_creation_tokens), 0) AS total_cache_creation,
COALESCE(SUM(cache_read_tokens), 0) AS total_cache_read,
COALESCE(SUM(output_tokens), 0) AS total_output,
COALESCE(SUM(message_count), 0) AS total_messages,
COUNT(DISTINCT session_id) AS total_sessions
FROM agent_token_usage
${tokenFilter.clause}`,
tokenFilter.params
);
return (
result.rows[0] ?? {
total_input: 0,
total_cache_creation: 0,
total_cache_read: 0,
total_output: 0,
total_messages: 0,
total_sessions: 0,
}
);
}

async function handleTokenDaily(
deps: ActivityRouteDeps,
request: FastifyRequest
) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const tokenFilter = deps.timeRangeClause(
aq,
"COALESCE(session_start, harvested_at)"
);
const result = await deps.pool.query<{
day: string;
input_tokens: number;
cache_creation_tokens: number;
cache_read_tokens: number;
output_tokens: number;
messages: number;
}>(
`SELECT
${deps.dateTruncTz(aq.granularity, "COALESCE(session_start, harvested_at)", aq.tz)} AS day,
SUM(input_tokens) AS input_tokens,
SUM(cache_creation_tokens) AS cache_creation_tokens,
SUM(cache_read_tokens) AS cache_read_tokens,
SUM(output_tokens) AS output_tokens,
SUM(message_count) AS messages
FROM agent_token_usage
${tokenFilter.clause}
GROUP BY day ORDER BY day`,
tokenFilter.params
);
return { days: result.rows, granularity: aq.granularity };
}

async function handleTokenByProject(
deps: ActivityRouteDeps,
request: FastifyRequest
) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const tokenFilter = deps.timeRangeClause(
aq,
"COALESCE(t.session_start, t.harvested_at)"
);
const result = await deps.pool.query<{
project_dir: string;
total_input: number;
total_output: number;
messages: number;
}>(
`SELECT
COALESCE(a.git_context->>'repoRoot', a.cwd) AS project_dir,
SUM(t.input_tokens + t.cache_creation_tokens + t.cache_read_tokens) AS total_input,
SUM(t.output_tokens) AS total_output,
SUM(t.message_count) AS messages
FROM agent_token_usage t
JOIN agents a ON a.id = t.agent_id
${tokenFilter.clause}
GROUP BY project_dir
ORDER BY total_input DESC
LIMIT 20`,
tokenFilter.params
);
return { projects: result.rows };
}

async function handleTokenByModel(
deps: ActivityRouteDeps,
request: FastifyRequest
) {
const aq = deps.parseActivityQuery(request.query as Record<string, unknown>);
const tokenFilter = deps.timeRangeClause(
aq,
"COALESCE(session_start, harvested_at)"
);
const result = await deps.pool.query<{
model: string;
total_input: number;
total_cache_creation: number;
total_cache_read: number;
total_output: number;
sessions: number;
}>(
`SELECT
model,
COALESCE(SUM(input_tokens), 0) AS total_input,
COALESCE(SUM(cache_creation_tokens), 0) AS total_cache_creation,
COALESCE(SUM(cache_read_tokens), 0) AS total_cache_read,
COALESCE(SUM(output_tokens), 0) AS total_output,
COUNT(DISTINCT session_id) AS sessions
FROM agent_token_usage
${tokenFilter.clause}
GROUP BY model
ORDER BY (SUM(input_tokens) + SUM(cache_creation_tokens) + SUM(cache_read_tokens) + SUM(output_tokens)) DESC`,
tokenFilter.params
);
return { models: result.rows };
}

async function handleHarvestTokens(
deps: ActivityRouteDeps,
request: FastifyRequest,
reply: FastifyReply
) {
const { id } = request.params as { id: string };
const agent = await deps.agentManager.getAgent(id);
if (!agent) {
return reply.code(404).send({ error: "Agent not found" });
}
await deps.agentManager.harvestAgentTokens(agent);
return { ok: true };
}
} from "../../activity-metrics.js";
import type { ActivityRouteDeps } from "./shared.js";

async function handleHistoryProjects(
deps: ActivityRouteDeps,
Expand Down Expand Up @@ -740,35 +429,10 @@ async function handleHistoryAgentDetail(
};
}

export async function registerActivityRoutes(
export async function registerActivityHistoryRoutes(
app: FastifyInstance,
deps: ActivityRouteDeps
): Promise<void> {
app.get("/api/v1/activity/heatmap", (req) => handleHeatmap(deps, req));
app.get("/api/v1/activity/stats", (req) => handleStats(deps, req));
app.get("/api/v1/activity/daily-status", (req) =>
handleDailyStatus(deps, req)
);
app.get("/api/v1/activity/active-hours", (req) =>
handleActiveHours(deps, req)
);
app.get("/api/v1/activity/agents-created", (req) =>
handleAgentsCreated(deps, req)
);
app.get("/api/v1/activity/working-time-by-project", (req) =>
handleWorkingTimeByProject(deps, req)
);
app.get("/api/v1/activity/token-stats", (req) => handleTokenStats(deps, req));
app.get("/api/v1/activity/token-daily", (req) => handleTokenDaily(deps, req));
app.get("/api/v1/activity/token-by-project", (req) =>
handleTokenByProject(deps, req)
);
app.get("/api/v1/activity/token-by-model", (req) =>
handleTokenByModel(deps, req)
);
app.post("/api/v1/agents/:id/harvest-tokens", (req, reply) =>
handleHarvestTokens(deps, req, reply)
);
app.get("/api/v1/history/projects", (req) =>
handleHistoryProjects(deps, req)
);
Expand Down
17 changes: 17 additions & 0 deletions apps/server/src/routes/activity/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
import type { FastifyInstance } from "fastify";

import type { ActivityRouteDeps } from "./shared.js";
import { registerActivityHistoryRoutes } from "./history-routes.js";
import { registerActivityMetricsRoutes } from "./metrics-routes.js";
import { registerActivityTokenRoutes } from "./token-routes.js";

export type { ActivityRouteDeps } from "./shared.js";

export async function registerActivityRoutes(
app: FastifyInstance,
deps: ActivityRouteDeps
): Promise<void> {
await registerActivityMetricsRoutes(app, deps);
await registerActivityTokenRoutes(app, deps);
await registerActivityHistoryRoutes(app, deps);
}
Loading
Loading