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
Expand Up @@ -68,6 +68,10 @@ enum class WorkThen(val raw: String, val label: String) {
data class WorkWhere(
val target: Target,
val detail: String?,
/** On a machine: the `local_hosts` id it runs on (the web's `where.hostId`). */
val hostId: String? = null,
/** On a machine: the directory it runs in, as stored. */
val dir: String? = null,
) {
enum class Target { POD, MACHINE }

Expand Down Expand Up @@ -454,7 +458,7 @@ object WorkFeed {
dir: String?,
): WorkWhere {
val parts = listOfNotNull(hostName[hostId.orEmpty()], shortDir(dir)).filter { it.isNotEmpty() }
return WorkWhere(WorkWhere.Target.MACHINE, parts.takeIf { it.isNotEmpty() }?.joinToString(" · "))
return WorkWhere(WorkWhere.Target.MACHINE, parts.takeIf { it.isNotEmpty() }?.joinToString(" · "), hostId, dir)
}

fun pod(detail: String?) = WorkWhere(WorkWhere.Target.POD, detail)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,7 @@ class WorkFeedTest {
assertEquals(listOf(WorkStatus.NEEDS_YOU, WorkStatus.NEEDS_YOU), rows.take(2).map { it.status })
assertEquals("terminal-lt1", rows.first().key, "most recent needs-you first")

assertEquals(WorkWhere(WorkWhere.Target.MACHINE, "M1 · ~/app"), rows.row("task-t2").where)
assertEquals(WorkWhere(WorkWhere.Target.MACHINE, "M1 · ~/app", "h1", "/Users/dev/app"), rows.row("task-t2").where)
assertEquals("PR 7", rows.row("task-t1").note)
assertEquals("/tasks/t1", rows.row("task-t1").href)
assertEquals(WorkStatus.PAUSED, rows.row("job-j1").status)
Expand Down Expand Up @@ -284,7 +284,7 @@ class WorkFeedTest {
assertNull(it.prUrl, "an empty PR link is no link")
}
rows.row("job-j").let {
assertEquals(WorkWhere(WorkWhere.Target.MACHINE, "/srv/jobs"), it.where, "an unknown host drops out")
assertEquals(WorkWhere(WorkWhere.Target.MACHINE, "/srv/jobs", "h9", "/srv/jobs"), it.where, "an unknown host drops out of the label")
assertEquals(WhenKind.TRIGGER, it.whenKind)
assertEquals("Claude Code", it.whoLabel)
}
Expand Down
30 changes: 22 additions & 8 deletions apps/api/src/routes/local.ts
Original file line number Diff line number Diff line change
Expand Up @@ -589,16 +589,21 @@ export async function localRoutes(rawApp: FastifyInstance) {
operationId: "getLocalTerminalTranscript",
summary: "The conversation of an agent session (prompts, replies, tool calls)",
description:
"Distilled by the daemon from the agent CLI's own transcript, so it covers the whole session — not just the last screen — and reads on any device. Grows while the session runs; `after` fetches only entries past a seq.",
"Distilled by the daemon from the agent CLI's own transcript, so it covers the whole session — not just the last screen — and reads on any device. Grows while the session runs; `after` fetches only entries past a seq. `before` instead returns the last `limit` entries before a seq (pass one past the highest seq for the latest page), with `hasEarlier` saying whether more precede them — how a phone opens a long session without downloading every tool output.",
tags: ["Local"],
params: z.object({ id: z.string().uuid() }),
querystring: z.object({
after: z.coerce.number().int().min(0).default(0),
before: z.coerce.number().int().min(1).optional(),
limit: z.coerce.number().int().min(1).max(5000).default(2000),
}),
response: {
200: z.object({
entries: z.array(LocalTranscriptEntrySchema),
hasEarlier: z
.boolean()
.optional()
.describe("With `before`: entries precede the first one returned"),
/** True when fewer than `limit` entries came back, i.e. the caller has everything stored. */
complete: z.boolean(),
backfilling: z
Expand All @@ -616,18 +621,27 @@ export async function localRoutes(rawApp: FastifyInstance) {
if (!terminal || !terminalService.canAccessTerminal(terminal, req.user?.id)) {
return reply.status(404).send({ error: "Terminal not found" });
}
const entries = await terminalService.getTranscript(
terminal.id,
req.query.after,
req.query.limit,
);
const { after, before, limit } = req.query;
const page =
before !== undefined
? await terminalService.getTranscriptBefore(terminal.id, before, limit)
: {
entries: await terminalService.getTranscript(terminal.id, after, limit),
hasEarlier: undefined,
};
const { entries } = page;
// Nothing stored at all: a finished session's machine can still read
// it off disk.
const backfilling =
req.query.after === 0 &&
(before !== undefined ? !page.hasEarlier : after === 0) &&
entries.length === 0 &&
terminalService.requestTranscriptBackfill(terminal);
reply.send({ entries, complete: entries.length < req.query.limit, backfilling });
reply.send({
entries,
complete: entries.length < limit,
backfilling,
...(page.hasEarlier !== undefined && { hasEarlier: page.hasEarlier }),
});
},
);

Expand Down
7 changes: 7 additions & 0 deletions apps/api/src/schemas/work.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,13 @@ export const WorkRowSchema = z
where: z.object({
target: z.enum(["pod", "machine"]),
detail: z.string().nullable(),
hostId: z.string().nullable().optional().describe("On a machine: the local host it runs on"),
hostName: z
.string()
.nullable()
.optional()
.describe("On a machine: its name, when it is one of the caller's machines"),
dir: z.string().nullable().optional().describe("On a machine: the directory it runs in"),
}),
who: z.string().describe('Agent runtime id, or "terminal"'),
then: z.enum(WORK_THENS),
Expand Down
58 changes: 44 additions & 14 deletions apps/api/src/services/local-terminal-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1067,20 +1067,50 @@ export async function getTranscript(
)
.orderBy(asc(localTerminalTranscripts.seq))
.limit(limit);
return rows.map((r) =>
reclassifyStoredTurn({
seq: r.seq,
role: r.role as LocalTranscriptRole,
kind: r.kind as LocalTranscriptKind,
source: (r.source as LocalTranscriptSource | null) ?? null,
text: r.text,
detail: r.detail,
toolName: r.toolName,
toolUseId: r.toolUseId,
isError: r.isError,
at: r.at ? r.at.toISOString() : null,
}),
);
return rows.map(transcriptRowToEntry);
}

/**
* The last `limit` entries before `beforeSeq`, in order — how a phone opens
* a long conversation: the latest exchange first, earlier pages on request,
* instead of every tool output since the session began. `hasEarlier` says
* whether anything precedes the first entry returned.
*/
export async function getTranscriptBefore(
terminalId: string,
beforeSeq: number,
limit: number,
): Promise<{ entries: LocalTranscriptEntry[]; hasEarlier: boolean }> {
const rows = await db
.select()
.from(localTerminalTranscripts)
.where(
and(
eq(localTerminalTranscripts.terminalId, terminalId),
lt(localTerminalTranscripts.seq, beforeSeq),
),
)
.orderBy(desc(localTerminalTranscripts.seq))
.limit(limit + 1);
const hasEarlier = rows.length > limit;
return { entries: rows.slice(0, limit).reverse().map(transcriptRowToEntry), hasEarlier };
}

function transcriptRowToEntry(
r: typeof localTerminalTranscripts.$inferSelect,
): LocalTranscriptEntry {
return reclassifyStoredTurn({
seq: r.seq,
role: r.role as LocalTranscriptRole,
kind: r.kind as LocalTranscriptKind,
source: (r.source as LocalTranscriptSource | null) ?? null,
text: r.text,
detail: r.detail,
toolName: r.toolName,
toolUseId: r.toolUseId,
isError: r.isError,
at: r.at ? r.at.toISOString() : null,
});
}

/**
Expand Down
16 changes: 16 additions & 0 deletions apps/api/src/services/local-transcript.int.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import {
deleteTerminal,
getTerminal,
getTranscript,
getTranscriptBefore,
handleExit,
handleSession,
handleStarted,
Expand Down Expand Up @@ -156,6 +157,21 @@ describe("handleTranscript / getTranscript", () => {
expect(await countTranscript(terminal.id)).toBe(4);
});

it("pages backwards from the latest entry", async () => {
const { host, terminal } = await runningAgentTerminal();
await handleTranscript(
host.id,
terminal.id,
[1, 2, 3, 4, 5].map((n) => entry(n)),
);
const latest = await getTranscriptBefore(terminal.id, 20_001, 2);
expect(latest.entries.map((e) => e.seq)).toEqual([4, 5]);
expect(latest.hasEarlier).toBe(true);
const earlier = await getTranscriptBefore(terminal.id, 4, 3);
expect(earlier.entries.map((e) => e.seq)).toEqual([1, 2, 3]);
expect(earlier.hasEarlier).toBe(false);
});

it("refuses writes from another host and after the terminal exited", async () => {
const { host, terminal } = await runningAgentTerminal();
const other = await makeHost();
Expand Down
7 changes: 6 additions & 1 deletion apps/api/src/services/work-service.int.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,12 @@ describe("listWork", () => {
expect(terminal).toBeLessThan(live);
expect(rows[terminal]).toMatchObject({
status: "needs_you",
where: { target: "machine", detail: "me-mac · ~/app" },
where: {
target: "machine",
detail: "me-mac · ~/app",
hostId: me.host.id,
hostName: "me-mac",
},
});
});
});
Expand Down
17 changes: 15 additions & 2 deletions apps/api/src/services/work-service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,13 @@ describe("projectWork", () => {
});

const t2 = rows.find((r) => r.key === "task-t2")!;
expect(t2.where).toEqual({ target: "machine", detail: "M1 · ~/app" });
expect(t2.where).toEqual({
target: "machine",
detail: "M1 · ~/app",
hostId: "h1",
hostName: "M1",
dir: "/Users/dev/app",
});
expect(t2).toMatchObject({ when: "on a trigger", spawned: true, then: "until-merged" });
expect(rows.find((r) => r.key === "task-t1")).toMatchObject({
note: "PR 7",
Expand Down Expand Up @@ -250,7 +256,14 @@ describe("projectWork", () => {
}),
);
expect(row).toMatchObject({
where: { target: "machine", detail: "~/app" },
// Another person's machine goes unnamed, but its id still groups the row.
where: {
target: "machine",
detail: "~/app",
hostId: "someone-elses",
hostName: null,
dir: "/home/dev/app",
},
who: "codex",
then: "exits",
when: "trigger",
Expand Down
17 changes: 12 additions & 5 deletions apps/api/src/services/work-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -215,11 +215,18 @@ function context(hosts: WorkSources["hosts"], triggers: TriggerRow[]): Context {
byId.set(t.id, t);
}
return {
machine: (hostId, dir) => ({
target: "machine",
detail:
[hostName.get(hostId ?? "") ?? null, shortDir(dir)].filter(Boolean).join(" · ") || null,
}),
machine: (hostId, dir) => {
// Only the caller's own machines are named; the id is kept either way
// so the Machines page can group work by the machine it runs on.
const name = hostName.get(hostId ?? "") ?? null;
return {
target: "machine",
detail: [name, shortDir(dir)].filter(Boolean).join(" · ") || null,
hostId: hostId ?? null,
hostName: name,
dir: dir ?? null,
};
},
triggersOf: (id) => distinctTriggers(byTarget.get(id) ?? []),
startedBy: (ticketSource, triggerId) => {
if (ticketSource) return [{ type: "ticket", source: ticketSource }];
Expand Down
Loading
Loading