Commit

Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request

A workspace's Usage page shows spend over a period, runs, the average run, credit left with how long it lasts at this rate, spend per day by kind of work, and breakdowns by task, repository, model and the pull requests that cost most. The sidebar shows this month's usage instead of only the balance. Also: people can message the agent at work on a pull request; the sandbox hands messages to the agent between its steps.

syntaqxcommitted Parent8e2c461Browse files
20 files+942−110/20 viewed
+36−1
5252 GetRepoSettings,
5353 UpdateRepoSettings,
5454 GetMergeQueue,
55+ MessageAgent,
56+ TakeMessages,
5557 ListIssues,
5658 GetIssue,
5759 CreateIssue,
185187 }
186188
187189 impl Op {
188− pub const ALL: [Op; 32] = [
190+ pub const ALL: [Op; 34] = [
189191 Op::Whoami,
190192 Op::CreateWorkspace,
191193 Op::ListRepos,
195197 Op::GetRepoSettings,
196198 Op::UpdateRepoSettings,
197199 Op::GetMergeQueue,
200+ Op::MessageAgent,
201+ Op::TakeMessages,
198202 Op::ListIssues,
199203 Op::GetIssue,
200204 Op::CreateIssue,
235239 Op::UpdateRepo => "update_repo",
236240 Op::GetRepoSettings => "get_repo_settings",
237241 Op::GetMergeQueue => "get_merge_queue",
242+ Op::MessageAgent => "message_agent",
243+ Op::TakeMessages => "take_messages",
238244 Op::UpdateRepoSettings => "update_repo_settings",
239245 Op::ListIssues => "list_issues",
240246 Op::GetIssue => "get_issue",
281287 Op::UpdateRepoSettings => {
282288 "Change how a repository handles pull requests. Only the fields given are changed. Members of its workspace only."
283289 }
290+ Op::MessageAgent => {
291+ "Send the agent working on a pull request a message: a correction, a hint, a change of plan. It receives it at its next step, and it is recorded in the pull request's session. The pull request's author and members of its workspace only."
292+ }
293+ Op::TakeMessages => {
294+ "For a g1t agent at work: the messages people have sent it that it has not seen yet. Each is returned once."
295+ }
284296 Op::GetMergeQueue => {
285297 "A repository's merge queue: the pull requests waiting to land, in order, each with the state it is being tested in (the default branch with the pull requests ahead of it merged in) and how that went; then those that recently landed or left. With the queue on, merging a pull request adds it here."
286298 }
387399 ),
388400 Op::GetRepoSettings => object(json!({ "repo": repo_schema() }), &["repo"]),
389401 Op::GetMergeQueue => object(json!({ "repo": repo_schema() }), &["repo"]),
402+ Op::MessageAgent => object(
403+ numbered(json!({
404+ "body": { "type": "string", "description": "What to tell the agent." },
405+ })),
406+ &["repo", "number", "body"],
407+ ),
408+ Op::TakeMessages => object(numbered(json!({})), &["repo", "number"]),
390409 Op::UpdateRepoSettings => object(
391410 json!({
392411 "repo": repo_schema(),
814833 Op::GetMergeQueue => {
815834 pass(work, "queue", &json!({ "repo": repo, "viewer": viewer })).await
816835 }
836+ Op::MessageAgent => {
837+ pass(
838+ work,
839+ "message_agent",
840+ &json!({ "actor": actor(), "repo": repo, "number": number, "body": text(input, "body") }),
841+ )
842+ .await
843+ }
844+ Op::TakeMessages => {
845+ pass(
846+ work,
847+ "take_messages",
848+ &json!({ "actor": actor(), "repo": repo, "number": number }),
849+ )
850+ .await
851+ }
817852 Op::UpdateRepoSettings => {
818853 // What is not given stays as it is.
819854 let current: Outcome<RepoSettings> = g1t_kit::call(
+12−0
4848 ),
4949 route("GET", "/repos/:owner/:name/queue", Op::GetMergeQueue, &[]),
5050 route(
51+ "POST",
52+ "/repos/:owner/:name/pulls/:number/messages",
53+ Op::MessageAgent,
54+ &[],
55+ ),
56+ route(
57+ "POST",
58+ "/repos/:owner/:name/pulls/:number/messages/take",
59+ Op::TakeMessages,
60+ &[],
61+ ),
62+ route(
5163 "GET",
5264 "/repos/:owner/:name/events",
5365 Op::ListEvents,
+23−9
11 import {
2+ BarChart3,
23 BookMarked,
34 BookOpen,
45 Check,
5657 } | null;
5758 /** The workspace's agent credit, if billing is on and they may see it. */
5859 creditMicros: number | null;
60+ /** What its agents have cost since the start of the month. */
61+ monthSpentMicros: number | null;
5962 };
6063
6164 function SidebarLink({
197200 <SidebarLink to={`/${ws.slug}/-/tokens`} icon={<KeyRound size={15} />}>
198201 Access tokens
199202 </SidebarLink>
203+ <SidebarLink to={`/${ws.slug}/-/usage`} icon={<BarChart3 size={15} />}>
204+ Usage
205+ </SidebarLink>
200206 <SidebarLink to={`/${ws.slug}/-/billing`} icon={<CreditCard size={15} />}>
201207 Billing
202208 </SidebarLink>
292298 </nav>
293299
294300 <div className="border-t border-line p-2">
295− {ws && shell.creditMicros != null && (
301+ {ws && shell.monthSpentMicros != null && (
296302 <Link
297− to={`/${ws.slug}/-/billing`}
298− className={`mb-1 flex items-center gap-2.5 rounded-md px-2 py-1.5 text-sm transition-colors hover:bg-raised ${
299− shell.creditMicros <= 0 ? "text-warn" : "text-muted"
300− }`}
303+ to={`/${ws.slug}/-/usage`}
304+ className="mb-1 block rounded-md px-2 py-1.5 text-sm text-muted transition-colors hover:bg-raised hover:text-fg"
301305 >
302− <CreditCard size={15} className="shrink-0" />
303− <span className="grow">Agent credit</span>
304− <span className="font-mono text-xs tabular-nums">
305− ${(shell.creditMicros / MICROS_PER_DOLLAR).toFixed(2)}
306+ <span className="flex items-center gap-2.5">
307+ <BarChart3 size={15} className="shrink-0 text-faint" />
308+ <span className="grow">Usage this month</span>
309+ <span className="font-mono text-xs text-fg tabular-nums">
310+ ${(shell.monthSpentMicros / MICROS_PER_DOLLAR).toFixed(2)}
311+ </span>
306312 </span>
313+ {shell.creditMicros != null && (
314+ <span
315+ className={`mt-0.5 block pl-6.5 text-xs ${shell.creditMicros <= 0 ? "text-warn" : "text-faint"}`}
316+ >
317+ ${(shell.creditMicros / MICROS_PER_DOLLAR).toFixed(2)} of credit left
318+ </span>
319+ )}
307320 </Link>
308321 )}
309322 <a
349362 for (const membership of user.workspaces ?? []) {
350363 commands.push(
351364 { label: membership.slug, hint: "Workspace", to: `/${membership.slug}`, icon: <Avatar name={membership.slug} size={15} square /> },
365+ { label: "Usage", hint: membership.slug, to: `/${membership.slug}/-/usage`, icon: <BarChart3 size={15} /> },
352366 { label: "Billing", hint: membership.slug, to: `/${membership.slug}/-/billing`, icon: <CreditCard size={15} /> },
353367 { label: "Access tokens", hint: membership.slug, to: `/${membership.slug}/-/tokens`, icon: <KeyRound size={15} /> },
354368 );
+5−1
8787 const here = params.owner ? memberships.find((m) => m.slug === params.owner?.toLowerCase()) : undefined;
8888 const workspace = here ?? memberships[0] ?? null;
8989 const path = params.owner && params.repo ? { namespace: params.owner, name: params.repo } : null;
90− const [listed, counts, account] = await Promise.all([
90+ const now = new Date();
91+ const monthStart = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), 1)).toISOString();
92+ const [listed, counts, account, usage] = await Promise.all([
9193 workspace ? repos.list(user, { namespace: workspace.slug }) : Promise.resolve([]),
9294 path ? work.counts(path, user) : Promise.resolve(null),
9395 workspace ? billing.account(workspace.slug, user) : Promise.resolve(null),
96+ workspace ? billing.usage(workspace.slug, user, monthStart) : Promise.resolve(null),
9497 ]);
9598 return {
9699 workspace,
108111 : null,
109112 creditMicros:
110113 account?.ok && account.value.status.enabled ? account.value.balanceMicros : null,
114+ monthSpentMicros: usage?.ok ? usage.value.spentMicros : null,
111115 };
112116 }
113117
+1−0
2020 index("routes/workspace/overview.tsx"),
2121 route("-/people", "routes/workspace/people.tsx"),
2222 route("-/tokens", "routes/workspace/tokens.tsx"),
23+ route("-/usage", "routes/workspace/usage.tsx"),
2324 route("-/billing", "routes/workspace/billing.tsx"),
2425 route("-/settings", "routes/workspace/settings.tsx"),
2526 ]),
+47−0
149149 })
150150 : action === "unqueue"
151151 ? await work.removeFromQueue(user, path, number)
152+ : action === "message"
153+ ? await work.messageAgent(user, path, number, String(form.get("body") ?? ""))
152154 : action === "close"
153155 ? await work.closePull(user, path, number)
154156 : action === "recheck"
315317 behind,
316318 reviewPending,
317319 lifecycle,
320+ messages,
318321 landing,
319322 stalled,
320323 requireUpToDate,
428431
429432 {lifecycle && <LifecyclePanel lifecycle={lifecycle} />}
430433
434+ {/* Steering: while its agent works, people can tell it things. */}
435+ {canManage &&
436+ pull.runtime === "hosted" &&
437+ (working || ["working", "revising", "catching_up"].includes(lifecycle?.stage ?? "")) && (
438+ <Form method="post" className="mt-4 rounded-2xl bg-surface p-4 ring-1 ring-merged/30">
439+ <p className="flex items-center gap-2 text-sm font-medium">
440+ <Sparkles size={15} className="text-merged" />
441+ Message the agent
442+ </p>
443+ <p className="mt-1 text-xs text-muted">
444+ A correction, a hint, a change of plan. It reads it at its next step, without
445+ starting over.
446+ </p>
447+ <div className="mt-3 flex gap-2">
448+ <input type="hidden" name="action" value="message" />
449+ <input
450+ name="body"
451+ required
452+ autoComplete="off"
453+ data-1p-ignore
454+ placeholder="Keep the old flag working too…"
455+ className="h-9 min-w-0 grow rounded-md bg-bg px-3 text-sm ring-1 ring-line outline-none placeholder:text-faint focus:ring-merged/60"
456+ />
457+ <Button type="submit">Send</Button>
458+ </div>
459+ </Form>
460+ )}
461+ {messages.length > 0 && (
462+ <ul className="mt-3 space-y-1.5">
463+ {messages.map((message) => (
464+ <li key={message.id} className="flex items-start gap-2 text-sm">
465+ <Avatar name={message.author} size={18} />
466+ <span className="min-w-0 grow">
467+ <span className="font-medium">{message.author}</span>{" "}
468+ <span className="text-muted">to the agent:</span> {message.body}
469+ </span>
470+ <span className={`shrink-0 text-xs ${message.deliveredAt ? "text-accent" : "text-faint"}`}>
471+ {message.deliveredAt ? "read by the agent" : "waiting for its next step"}
472+ </span>
473+ </li>
474+ ))}
475+ </ul>
476+ )}
477+
431478 {issue && (
432479 <Link
433480 to={`${base}/issues/${issue.number}`}
+299−0
1+import { ArrowUpRight, CreditCard } from "lucide-react";
2+import { Link, data } from "react-router";
3+
4+import { MICROS_PER_DOLLAR, type UsageSlice } from "@g1t/contracts";
5+
6+import type { Route } from "./+types/usage";
7+import { ButtonLink } from "../../components/ui";
8+import { billing } from "../../lib/services.server";
9+import { getViewer, roleIn, unwrap } from "../../lib/session.server";
10+
11+const PERIODS = {
12+ month: "This month",
13+ "7d": "Last 7 days",
14+ "30d": "Last 30 days",
15+ "90d": "Last 90 days",
16+} as const;
17+type Period = keyof typeof PERIODS;
18+
19+/** What each kind of agent work is called, and its colour. */
20+const TASKS: Record<string, { label: string; color: string }> = {
21+ implement: { label: "Making changes", color: "var(--color-merged)" },
22+ review: { label: "Reviews", color: "var(--color-info)" },
23+ revise: { label: "Revisions", color: "var(--color-warn)" },
24+ update: { label: "Catching up", color: "var(--color-accent)" },
25+ plan: { label: "Planning", color: "#f0a6ca" },
26+ other: { label: "Other", color: "var(--color-faint)" },
27+};
28+const task = (key: string) => TASKS[key] ?? TASKS.other!;
29+
30+function start(period: Period, now = new Date()): Date {
31+ if (period === "month") return new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), 1));
32+ const days = Number(period.replace("d", ""));
33+ const day = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), now.getUTCDate()));
34+ return new Date(day.getTime() - (days - 1) * 86_400_000);
35+}
36+
37+export function meta({ params }: Route.MetaArgs) {
38+ return [{ title: `Usage · ${params.owner} · g1t` }];
39+}
40+
41+export async function loader({ params, context, request }: Route.LoaderArgs) {
42+ const viewer = getViewer(context);
43+ if (!roleIn(viewer, params.owner)) throw data(null, { status: 404 });
44+ const asked = new URL(request.url).searchParams.get("period");
45+ const period: Period = asked && asked in PERIODS ? (asked as Period) : "month";
46+ const since = start(period);
47+ const [usage, account] = await Promise.all([
48+ billing.usage(params.owner, viewer, since.toISOString()),
49+ billing.account(params.owner, viewer),
50+ ]);
51+ return { period, since: since.toISOString(), usage: unwrap(usage), account: unwrap(account) };
52+}
53+
54+function dollars(micros: number, digits = 2): string {
55+ return `$${(micros / MICROS_PER_DOLLAR).toFixed(digits)}`;
56+}
57+
58+function Stat({ label, value, note }: { label: string; value: string; note?: string }) {
59+ return (
60+ <div className="rounded-2xl bg-surface p-5 ring-1 ring-line">
61+ <p className="text-sm text-muted">{label}</p>
62+ <p className="mt-2 text-3xl font-semibold tracking-tight tabular-nums">{value}</p>
63+ {note && <p className="mt-1 text-xs text-faint">{note}</p>}
64+ </div>
65+ );
66+}
67+
68+/** Spend per day, stacked by kind of work. */
69+function Chart({ since, slices }: { since: string; slices: UsageSlice[] }) {
70+ const first = new Date(since);
71+ const today = new Date();
72+ const days: string[] = [];
73+ for (let day = first; day <= today; day = new Date(day.getTime() + 86_400_000)) {
74+ days.push(day.toISOString().slice(0, 10));
75+ }
76+ const byDay = new Map<string, { task: string; micros: number }[]>();
77+ for (const slice of slices) {
78+ const [day, kind] = slice.key.split("/");
79+ byDay.set(day!, [...(byDay.get(day!) ?? []), { task: kind ?? "other", micros: slice.micros }]);
80+ }
81+ const totals = days.map((day) => (byDay.get(day) ?? []).reduce((sum, part) => sum + part.micros, 0));
82+ const max = Math.max(...totals, 1);
83+ const height = 160;
84+ const width = 100 / days.length;
85+ return (
86+ <div>
87+ <div className="relative" style={{ height }}>
88+ {/* Gridlines at a quarter, a half and three quarters of the highest day. */}
89+ {[0.25, 0.5, 0.75, 1].map((at) => (
90+ <div key={at} className="absolute inset-x-0 border-t border-line/60" style={{ bottom: `${at * 100}%` }}>
91+ <span className="absolute -top-2.5 right-0 bg-surface pl-1 font-mono text-[0.625rem] text-faint">
92+ {dollars(max * at)}
93+ </span>
94+ </div>
95+ ))}
96+ <div className="absolute inset-0 flex items-end gap-px pr-12">
97+ {days.map((day, index) => {
98+ const parts = byDay.get(day) ?? [];
99+ return (
100+ <div
101+ key={day}
102+ className="group relative flex h-full flex-col justify-end"
103+ style={{ width: `${width}%` }}
104+ title={`${day}: ${dollars(totals[index]!)}`}
105+ >
106+ {parts.map((part) => (
107+ <div
108+ key={part.task}
109+ className="w-full first:rounded-t-sm"
110+ style={{ height: `${(part.micros / max) * 100}%`, background: task(part.task).color }}
111+ />
112+ ))}
113+ <div className="pointer-events-none absolute bottom-full left-1/2 z-10 mb-2 hidden -translate-x-1/2 rounded-md bg-raised px-2 py-1 text-xs whitespace-nowrap ring-1 ring-line-strong group-hover:block">
114+ <span className="text-muted">{day}</span> {dollars(totals[index]!)}
115+ </div>
116+ </div>
117+ );
118+ })}
119+ </div>
120+ </div>
121+ <div className="mt-2 flex justify-between pr-12 font-mono text-[0.625rem] text-faint">
122+ <span>{days[0]}</span>
123+ <span>{days[days.length - 1]}</span>
124+ </div>
125+ </div>
126+ );
127+}
128+
129+function Breakdown({
130+ title,
131+ slices,
132+ total,
133+ label,
134+ link,
135+ color,
136+}: {
137+ title: string;
138+ slices: UsageSlice[];
139+ total: number;
140+ label?: (key: string) => string;
141+ link?: (key: string) => string | null;
142+ color?: (key: string) => string;
143+}) {
144+ return (
145+ <section className="rounded-2xl bg-surface p-5 ring-1 ring-line">
146+ <h3 className="text-sm font-medium">{title}</h3>
147+ {slices.length === 0 ? (
148+ <p className="mt-4 text-sm text-faint">Nothing in this period.</p>
149+ ) : (
150+ <ul className="mt-4 space-y-3">
151+ {slices.map((slice) => {
152+ const to = link?.(slice.key);
153+ const name = label ? label(slice.key) : slice.key;
154+ return (
155+ <li key={slice.key}>
156+ <div className="flex items-baseline justify-between gap-3 text-sm">
157+ <span className="min-w-0 truncate">
158+ {to ? (
159+ <Link to={to} className="hover:underline">
160+ {name}
161+ </Link>
162+ ) : (
163+ name
164+ )}
165+ </span>
166+ <span className="shrink-0 font-mono text-xs tabular-nums">
167+ {dollars(slice.micros)} <span className="text-faint">· {slice.runs} {slice.runs === 1 ? "run" : "runs"}</span>
168+ </span>
169+ </div>
170+ <div className="mt-1.5 h-1.5 overflow-hidden rounded-full bg-raised">
171+ <div
172+ className="h-full rounded-full"
173+ style={{
174+ width: `${total ? (slice.micros / total) * 100 : 0}%`,
175+ background: color?.(slice.key) ?? "var(--color-merged)",
176+ }}
177+ />
178+ </div>
179+ </li>
180+ );
181+ })}
182+ </ul>
183+ )}
184+ </section>
185+ );
186+}
187+
188+export default function UsagePage({ loaderData, params }: Route.ComponentProps) {
189+ const { period, since, usage, account } = loaderData;
190+ const base = `/${params.owner}`;
191+ const days = Math.max(1, Math.ceil((Date.now() - new Date(since).getTime()) / 86_400_000));
192+ const perDay = usage.spentMicros / days;
193+ const runway = perDay > 0 ? Math.floor(account.balanceMicros / perDay) : null;
194+ return (
195+ <div className="space-y-8">
196+ <div className="flex flex-wrap items-end justify-between gap-4">
197+ <div>
198+ <h2 className="text-xl font-semibold tracking-tight">Usage</h2>
199+ <p className="mt-1 max-w-2xl text-sm text-muted">
200+ What this workspace's agents cost: every run, charged at what the model provider charges
201+ plus {account.marginPercent}%, from the workspace's credit.
202+ </p>
203+ </div>
204+ <nav className="flex rounded-lg bg-surface p-1 ring-1 ring-line">
205+ {(Object.keys(PERIODS) as Period[]).map((key) => (
206+ <Link
207+ key={key}
208+ to={`?period=${key}`}
209+ preventScrollReset
210+ className={`rounded-md px-3 py-1.5 text-sm transition-colors ${
211+ key === period ? "bg-raised text-fg" : "text-muted hover:text-fg"
212+ }`}
213+ >
214+ {PERIODS[key]}
215+ </Link>
216+ ))}
217+ </nav>
218+ </div>
219+
220+ <div className="grid gap-4 sm:grid-cols-2 lg:grid-cols-4">
221+ <Stat label="Spent" value={dollars(usage.spentMicros)} note={`${dollars(usage.costMicros)} of it the model provider's`} />
222+ <Stat label="Agent runs" value={String(usage.runs)} note={PERIODS[period]} />
223+ <Stat
224+ label="Average run"
225+ value={usage.runs ? dollars(usage.spentMicros / usage.runs, 3) : "—"}
226+ note="Making a change, reviewing, revising…"
227+ />
228+ <div className="rounded-2xl bg-surface p-5 ring-1 ring-line">
229+ <p className="flex items-center justify-between text-sm text-muted">
230+ Credit left
231+ <Link to={`${base}/-/billing`} className="text-xs text-accent hover:underline">
232+ Add credit
233+ </Link>
234+ </p>
235+ <p className={`mt-2 text-3xl font-semibold tracking-tight tabular-nums ${account.balanceMicros <= 0 ? "text-warn" : ""}`}>
236+ {dollars(account.balanceMicros)}
237+ </p>
238+ <p className="mt-1 text-xs text-faint">
239+ {runway == null ? "Nothing spent in this period." : `About ${runway} ${runway === 1 ? "day" : "days"} at this rate.`}
240+ </p>
241+ </div>
242+ </div>
243+
244+ <section className="rounded-2xl bg-surface p-5 ring-1 ring-line">
245+ <div className="flex flex-wrap items-baseline justify-between gap-3">
246+ <h3 className="text-sm font-medium">Spend per day</h3>
247+ <ul className="flex flex-wrap gap-x-4 gap-y-1 text-xs text-muted">
248+ {usage.byTask.map((slice) => (
249+ <li key={slice.key} className="flex items-center gap-1.5">
250+ <span className="size-2 rounded-sm" style={{ background: task(slice.key).color }} />
251+ {task(slice.key).label}
252+ </li>
253+ ))}
254+ </ul>
255+ </div>
256+ <div className="mt-6">
257+ <Chart since={since} slices={usage.byDay} />
258+ </div>
259+ </section>
260+
261+ <div className="grid gap-4 lg:grid-cols-2">
262+ <Breakdown
263+ title="By kind of work"
264+ slices={usage.byTask}
265+ total={usage.spentMicros}
266+ label={(key) => task(key).label}
267+ color={(key) => task(key).color}
268+ />
269+ <Breakdown title="By repository" slices={usage.byRepo} total={usage.spentMicros} link={(key) => `/${key}`} />
270+ <Breakdown
271+ title="Pull requests that cost most"
272+ slices={usage.byPull}
273+ total={usage.byPull[0]?.micros ?? 0}
274+ // Planning runs belong to a repository, not a pull request.
275+ label={(key) => (key.endsWith("#0") ? `${key.slice(0, -2)} · planning` : key)}
276+ link={(key) => {
277+ const [repo, number] = key.split("#");
278+ if (!repo) return null;
279+ return number && number !== "0" ? `/${repo}/pull/${number}` : `/${repo}/plans`;
280+ }}
281+ />
282+ <Breakdown title="By model" slices={usage.byModel} total={usage.spentMicros} color={() => "var(--color-info)"} />
283+ </div>
284+
285+ <div className="flex flex-wrap items-center justify-between gap-3 rounded-2xl bg-surface px-5 py-4 text-sm ring-1 ring-line">
286+ <span className="flex items-center gap-2 text-muted">
287+ <CreditCard size={15} />
288+ {usage.addedMicros > 0
289+ ? `${dollars(usage.addedMicros)} of credit added in this period.`
290+ : "Every charge and payment is on the workspace's statement."}
291+ </span>
292+ <ButtonLink to={`${base}/-/billing`} variant="quiet">
293+ Statement and credit
294+ <ArrowUpRight size={14} />
295+ </ButtonLink>
296+ </div>
297+ </div>
298+ );
299+}
+44−0
154154 #[serde(default)]
155155 pub turns: u32,
156156 }
157+
158+
159+/// `usage`: what a workspace's agents cost over a period, broken down.
160+/// Members only. Returns `Outcome<Usage>`.
161+#[derive(Debug, Serialize, Deserialize)]
162+pub struct UsageArgs {
163+ pub workspace: String,
164+ pub viewer: Viewer,
165+ /// RFC 3339: the start of the period. The period runs to now.
166+ pub since: String,
167+}
168+
169+/// One slice of usage: what it was for, what it cost, how many runs.
170+#[derive(Clone, Debug, Serialize, Deserialize)]
171+#[serde(rename_all = "camelCase")]
172+pub struct UsageSlice {
173+ pub key: String,
174+ pub micros: i64,
175+ pub runs: u32,
176+}
177+
178+/// What a workspace's agents cost over a period.
179+#[derive(Clone, Debug, Serialize, Deserialize)]
180+#[serde(rename_all = "camelCase")]
181+pub struct Usage {
182+ pub since: String,
183+ /// Charged, including g1t's margin.
184+ pub spent_micros: i64,
185+ /// What the model provider charged, before the margin.
186+ pub cost_micros: i64,
187+ pub runs: u32,
188+ /// Spend per day (`YYYY-MM-DD`) and task, as `day/task` keys.
189+ pub by_day: Vec<UsageSlice>,
190+ /// Per task: implement, review, revise, update, plan.
191+ pub by_task: Vec<UsageSlice>,
192+ /// Per repository, `namespace/name`.
193+ pub by_repo: Vec<UsageSlice>,
194+ /// The pull requests that cost most, as `namespace/name#number`.
195+ pub by_pull: Vec<UsageSlice>,
196+ /// Per model, by its public name.
197+ pub by_model: Vec<UsageSlice>,
198+ /// Credit bought in the period.
199+ pub added_micros: i64,
200+}
+38−0
412412 /// be completed, for example.
413413 #[serde(default)]
414414 pub stalled: Option<String>,
415+ /// Messages people sent the agent while it worked, oldest first.
416+ #[serde(default)]
417+ pub messages: Vec<AgentMessage>,
418+}
419+
420+/// A message a person sent an agent at work on a pull request. The agent
421+/// receives it at its next step.
422+#[derive(Clone, Debug, Serialize, Deserialize)]
423+#[serde(rename_all = "camelCase")]
424+pub struct AgentMessage {
425+ pub id: String,
426+ pub author: String,
427+ pub body: String,
428+ /// RFC 3339.
429+ pub created_at: String,
430+ /// RFC 3339. When the agent received it; null until then.
431+ pub delivered_at: Option<String>,
432+}
433+
434+/// `message_agent`: sends the agent working on a pull request a message.
435+/// The pull request's author and members of the workspace may. Returns
436+/// `Outcome<AgentMessage>`.
437+#[derive(Debug, Serialize, Deserialize)]
438+pub struct MessageAgentArgs {
439+ pub actor: User,
440+ pub repo: RepoPath,
441+ pub number: u32,
442+ pub body: String,
443+}
444+
445+/// `take_messages`: the messages not yet delivered to the agent working on
446+/// a pull request, marked delivered. Only g1t's agents may. Returns
447+/// `Outcome<Vec<AgentMessage>>`.
448+#[derive(Debug, Serialize, Deserialize)]
449+pub struct TakeMessagesArgs {
450+ pub actor: User,
451+ pub repo: RepoPath,
452+ pub number: u32,
415453 }
416454
417455 /// A step on the way from an assigned issue to a pull request that is ready
+25−0
146146 if std::fs::write(path, config.to_string()).is_ok() {
147147 tools = vec!["--mcp-config".to_owned(), path.to_owned()];
148148 }
149+ // People can message the agent while it works: after each tool call
150+ // a hook asks g1t for messages and hands any to the agent.
151+ if let (Ok(repo), Ok(number)) = (std::env::var("G1T_REPO"), std::env::var("PULL_NUMBER")) {
152+ let steer = serde_json::json!({
153+ "api": std::env::var("G1T_API").unwrap_or_else(|_| "https://api.g1t.sh".to_owned()),
154+ "token": token,
155+ "repo": repo,
156+ "number": number.parse::<u32>().unwrap_or_default(),
157+ });
158+ let hooks = serde_json::json!({
159+ "hooks": {
160+ "PostToolUse": [{
161+ "matcher": "*",
162+ "hooks": [{ "type": "command", "command": "MODE=steer /usr/local/bin/g1t-runner", "timeout": 15 }],
163+ }],
164+ }
165+ });
166+ let home = std::env::var("HOME").unwrap_or_else(|_| "/home/node".to_owned());
167+ let settings = format!("{home}/.claude/settings.json");
168+ if std::fs::write(crate::steer::CONFIG, steer.to_string()).is_ok()
169+ && std::fs::create_dir_all(format!("{home}/.claude")).is_ok()
170+ {
171+ let _ = std::fs::write(settings, hooks.to_string());
172+ }
173+ }
149174 }
150175 let mut child = Command::new("claude")
151176 .current_dir(workdir)
+2−0
2828 mod report;
2929 mod review;
3030 mod revise;
31+mod steer;
3132 mod update;
3233
3334 use std::path::Path;
179180 Ok("revise") => std::process::exit(revise::main()),
180181 Ok("plan") => std::process::exit(plan::main()),
181182 Ok("queue") => std::process::exit(queue::main()),
183+ Ok("steer") => std::process::exit(steer::main()),
182184 _ => {}
183185 }
184186 let mut reporter = match Reporter::from_env() {
+86−0
1+//! Delivers people's messages to the agent while it works.
2+//!
3+//! Claude Code runs this as a hook after each of the agent's tool calls
4+//! (see `harness`). It asks g1t for messages the agent has not seen and,
5+//! if there are any, hands them to the agent as context for its next step.
6+//! It never fails the agent's run: any problem means no message this time.
7+//!
8+//! What it needs is in `/work/g1t-steer.json`, written by the harness:
9+//! the API, the agent's token, the repository and the pull request.
10+
11+use std::time::{SystemTime, UNIX_EPOCH};
12+
13+use serde::Deserialize;
14+
15+/// Where the harness leaves what this needs.
16+pub const CONFIG: &str = "/work/g1t-steer.json";
17+/// When it last asked, so that a burst of tool calls asks once.
18+const LAST_ASKED: &str = "/work/.g1t-steer-at";
19+/// How long to wait between asks.
20+const INTERVAL_MS: u128 = 10_000;
21+
22+#[derive(Deserialize)]
23+struct Config {
24+ api: String,
25+ token: String,
26+ repo: String,
27+ number: u32,
28+}
29+
30+#[derive(Deserialize)]
31+struct Message {
32+ author: String,
33+ body: String,
34+}
35+
36+fn now_ms() -> u128 {
37+ SystemTime::now()
38+ .duration_since(UNIX_EPOCH)
39+ .map(|elapsed| elapsed.as_millis())
40+ .unwrap_or_default()
41+}
42+
43+fn take() -> Option<Vec<Message>> {
44+ let config: Config = serde_json::from_str(&std::fs::read_to_string(CONFIG).ok()?).ok()?;
45+ let last: u128 = std::fs::read_to_string(LAST_ASKED)
46+ .ok()
47+ .and_then(|text| text.trim().parse().ok())
48+ .unwrap_or_default();
49+ let now = now_ms();
50+ if now.saturating_sub(last) < INTERVAL_MS {
51+ return None;
52+ }
53+ let _ = std::fs::write(LAST_ASKED, now.to_string());
54+ let response = ureq::post(&format!(
55+ "{}/repos/{}/pulls/{}/messages/take",
56+ config.api, config.repo, config.number
57+ ))
58+ .set("Authorization", &format!("Bearer {}", config.token))
59+ .send_json(serde_json::json!({}))
60+ .ok()?;
61+ response.into_json().ok()
62+}
63+
64+pub fn main() -> i32 {
65+ let Some(messages) = take().filter(|messages| !messages.is_empty()) else {
66+ return 0;
67+ };
68+ let said: Vec<String> = messages
69+ .iter()
70+ .map(|message| format!("{} says: {}", message.author, message.body))
71+ .collect();
72+ let context = format!(
73+ "A person watching your work just sent you a message on the pull request. Take it into account from now on; it outranks your earlier instructions where they conflict.\n\n{}",
74+ said.join("\n\n")
75+ );
76+ println!(
77+ "{}",
78+ serde_json::json!({
79+ "hookSpecificOutput": {
80+ "hookEventName": "PostToolUse",
81+ "additionalContext": context,
82+ }
83+ })
84+ );
85+ 0
86+}
+25−0
6464 account(workspace: string, viewer: Viewer): Promise<Result<BillingAccount>>;
6565 /** Newest first. Members of the workspace only. */
6666 ledger(workspace: string, viewer: Viewer): Promise<Result<LedgerEntry[]>>;
67+ /** What the workspace's agents cost since `since`, broken down. Members only. */
68+ usage(workspace: string, viewer: Viewer, since: string): Promise<Result<Usage>>;
6769 /**
6870 * Starts a card payment for credit and returns the page to send the
6971 * person to. Owners only. The payment's id comes back to `returnUrl` as
8991 model: string;
9092 }): Promise<Result<RunTicket | null>>;
9193 }
94+
95+
96+/** One slice of usage: what it was for, what it cost, how many runs. */
97+export type UsageSlice = { key: string; micros: number; runs: number };
98+
99+/** What a workspace's agents cost over a period. */
100+export type Usage = {
101+ since: string;
102+ /** Charged, including g1t's margin. */
103+ spentMicros: number;
104+ /** What the model provider charged, before the margin. */
105+ costMicros: number;
106+ runs: number;
107+ /** Spend per day and task, keyed `YYYY-MM-DD/task`. */
108+ byDay: UsageSlice[];
109+ byTask: UsageSlice[];
110+ byRepo: UsageSlice[];
111+ /** Keyed `namespace/name#number`. */
112+ byPull: UsageSlice[];
113+ byModel: UsageSlice[];
114+ /** Credit bought in the period. */
115+ addedMicros: number;
116+};
+2−0
138138 queueBuild: (repoId) => call("queue_build", { repoId }),
139139 failQueue: (entryId, token, error) => call("report_queue", { entryId, token, error }),
140140 removeFromQueue: (actor, repo, number) => call("remove_from_queue", { actor, repo, number }),
141+ messageAgent: (actor, repo, number, body) => call("message_agent", { actor, repo, number, body }),
141142 catchUpJob: (pullId) => call("catch_up_job", { pullId }),
142143 getSettings: (repo, viewer) => call("get_settings", { repo, viewer }),
143144 updateSettings: (actor, repo, settings) =>
176177 status: () => call("status", {}),
177178 account: (workspace, viewer) => call("account", { workspace, viewer }),
178179 ledger: (workspace, viewer) => call("ledger", { workspace, viewer }),
180+ usage: (workspace, viewer, since) => call("usage", { workspace, viewer, since }),
179181 checkout: (actor, workspace, amountCents, returnUrl) =>
180182 call("checkout", { actor, workspace, amountCents, returnUrl }),
181183 confirm: (workspace, viewer, session) => call("confirm", { workspace, viewer, session }),
+14−0
285285 landing: boolean;
286286 /** Why g1t stopped working on it, if it did. */
287287 stalled: string | null;
288+ /** Messages people sent the agent while it worked, oldest first. */
289+ messages: AgentMessage[];
288290 };
289291
290292 /**
560562 queueBuild(repoId: string): Promise<QueueJob[]>;
561563 /** Reports that a combined state could not be built or checked. */
562564 failQueue(entryId: string, token: string, error: string): Promise<Result<QueueState>>;
565+ /** Sends the agent working on a pull request a message, for its next step. */
566+ messageAgent(actor: User, repo: RepoPath, number: number, body: string): Promise<Result<AgentMessage>>;
563567 /** Takes a pull request out of the merge queue. Members only. */
564568 removeFromQueue(actor: User, repo: RepoPath, number: number): Promise<Result<Pull>>;
565569
696700 agent: string | null;
697701 };
698702
703+/** A message a person sent an agent at work on a pull request. */
704+export type AgentMessage = {
705+ id: string;
706+ author: string;
707+ body: string;
708+ createdAt: string;
709+ /** When the agent received it; null until then. */
710+ deliveredAt: string | null;
711+};
712+
699713 export type QueueState = "waiting" | "testing" | "passed" | "failed" | "landed" | "removed";
700714
701715 /** One pull request's place in a merge queue. */
+81−0
234234 ))
235235 }
236236
237+ async fn usage(&self, a: UsageArgs) -> Result<Outcome<Usage>> {
238+ let workspace = a.workspace.to_lowercase();
239+ if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
240+ return Ok(members_only());
241+ }
242+ #[derive(serde::Deserialize)]
243+ struct SliceRow {
244+ key: Option<String>,
245+ micros: Option<i64>,
246+ runs: Option<u32>,
247+ }
248+ let slices = |key: &str, limit: u32| {
249+ format!(
250+ "SELECT {key} AS key, -SUM(amount_micros) AS micros, COUNT(*) AS runs FROM ledger
251+ WHERE workspace = ?1 AND kind = 'usage' AND created_at >= ?2
252+ GROUP BY 1 ORDER BY micros DESC LIMIT {limit}"
253+ )
254+ };
255+ let query = |sql: String| {
256+ let db = &self.db;
257+ let workspace = workspace.clone();
258+ let since = a.since.clone();
259+ async move {
260+ let rows = db
261+ .prepare(sql)
262+ .bind(&[workspace.into(), since.into()])?
263+ .all()
264+ .await?
265+ .results::<SliceRow>()?;
266+ Ok::<Vec<UsageSlice>, worker::Error>(
267+ rows.into_iter()
268+ .map(|row| UsageSlice {
269+ key: row.key.unwrap_or_else(|| "other".to_owned()),
270+ micros: row.micros.unwrap_or_default(),
271+ runs: row.runs.unwrap_or_default(),
272+ })
273+ .collect(),
274+ )
275+ }
276+ };
277+ #[derive(serde::Deserialize)]
278+ struct Totals {
279+ spent: Option<i64>,
280+ cost: Option<i64>,
281+ runs: Option<u32>,
282+ added: Option<i64>,
283+ }
284+ let totals = self
285+ .db
286+ .prepare(
287+ "SELECT
288+ -SUM(CASE WHEN kind = 'usage' THEN amount_micros END) AS spent,
289+ SUM(CASE WHEN kind = 'usage' THEN cost_micros END) AS cost,
290+ SUM(CASE WHEN kind = 'usage' THEN 1 ELSE 0 END) AS runs,
291+ SUM(CASE WHEN kind = 'top_up' THEN amount_micros END) AS added
292+ FROM ledger WHERE workspace = ?1 AND created_at >= ?2",
293+ )
294+ .bind(&[workspace.as_str().into(), a.since.as_str().into()])?
295+ .first::<Totals>(None)
296+ .await?;
297+ let totals = totals.unwrap_or(Totals {
298+ spent: None,
299+ cost: None,
300+ runs: None,
301+ added: None,
302+ });
303+ Ok(Outcome::Ok(Usage {
304+ spent_micros: totals.spent.unwrap_or_default(),
305+ cost_micros: totals.cost.unwrap_or_default(),
306+ runs: totals.runs.unwrap_or_default(),
307+ added_micros: totals.added.unwrap_or_default(),
308+ by_day: query(slices("substr(created_at, 1, 10) || '/' || COALESCE(task, 'other')", 400)).await?,
309+ by_task: query(slices("task", 20)).await?,
310+ by_repo: query(slices("repo", 20)).await?,
311+ by_pull: query(slices("repo || '#' || number", 10)).await?,
312+ by_model: query(slices("model", 10)).await?,
313+ since: a.since,
314+ }))
315+ }
316+
237317 async fn checkout(&self, a: CheckoutArgs) -> Result<Outcome<Checkout>> {
238318 let workspace = a.workspace.to_lowercase();
239319 if a.actor.role_in(&workspace) != Some(Role::Owner) {
485565 "status" => reply(&billing.status()),
486566 "account" => reply(&billing.account(args(body)?).await?),
487567 "ledger" => reply(&billing.ledger(args(body)?).await?),
568+ "usage" => reply(&billing.usage(args(body)?).await?),
488569 "checkout" => reply(&billing.checkout(args(body)?).await?),
489570 "confirm" => reply(&billing.confirm(args(body)?).await?),
490571 "can_start" => reply(&billing.can_start(args(body)?).await?),
+2−0
207207 "read_session",
208208 "get_merge_queue",
209209 "list_events",
210+ // Messages people send it while it works, picked up between steps.
211+ "take_messages",
210212 ];
211213
212214 /** How an agent is told to use g1t's tools to work with the others. */
+12−0
1+-- Messages people send an agent while it works on a pull request. The
2+-- sandbox takes the undelivered ones at each step of the agent's run.
3+CREATE TABLE agent_messages (
4+ id TEXT PRIMARY KEY,
5+ pull_id TEXT NOT NULL,
6+ author_id TEXT NOT NULL,
7+ author_name TEXT NOT NULL,
8+ body TEXT NOT NULL,
9+ created_at TEXT NOT NULL,
10+ delivered_at TEXT
11+);
12+CREATE INDEX agent_messages_by_pull ON agent_messages (pull_id, created_at);
+34−0
77 mod checks;
88 mod lifecycle;
99 mod plans;
10+mod messages;
1011 mod queue;
1112 mod reviews;
1213 mod rows;
10611062 lifecycle,
10621063 landing,
10631064 stalled,
1065+ messages: self.messages(&pull.id).await?,
10641066 issue,
10651067 pull,
10661068 }))
15821584 Ok(Outcome::Ok(Appended { count }))
15831585 }
15841586
1587+ /// Adds entries to a pull request's session, each taking the next
1588+ /// sequence number, without announcing it.
1589+ pub(crate) async fn append_entries(&self, pull: &Pull, entries: &[NewSessionEntry]) -> Result<()> {
1590+ let now = rfc3339(now_ms());
1591+ let mut statements = Vec::with_capacity(entries.len());
1592+ for entry in entries {
1593+ let kind = serde_json::to_value(entry.kind)?;
1594+ let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1595+ statements.push(
1596+ self.db
1597+ .prepare(
1598+ "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1599+ SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1600+ FROM session_entries WHERE pull_id = ?",
1601+ )
1602+ .bind(&[
1603+ pull.id.as_str().into(),
1604+ kind.as_str().unwrap_or("note").into(),
1605+ text.into(),
1606+ optional(&entry.tool),
1607+ optional(&entry.commit.clone().or_else(|| pull.head_commit.clone())),
1608+ now.as_str().into(),
1609+ pull.id.as_str().into(),
1610+ ])?,
1611+ );
1612+ }
1613+ self.db.batch(statements).await?;
1614+ Ok(())
1615+ }
1616+
15851617 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
15861618 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
15871619 let rows = self
17281760 "queue_build" => reply(&work.queue_build(args(body)?).await?),
17291761 "report_queue" => reply(&work.report_queue(args(body)?).await?),
17301762 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
1763+ "message_agent" => reply(&work.message_agent(args(body)?).await?),
1764+ "take_messages" => reply(&work.take_messages(args(body)?).await?),
17311765 "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?),
17321766 "get_settings" => reply(&work.get_settings(args(body)?).await?),
17331767 "update_settings" => reply(&work.update_settings(args(body)?).await?),
+154−0
1+//! Steering: people send the agent at work on a pull request a message,
2+//! and the agent receives it at its next step. The sandbox asks for
3+//! undelivered messages after each of the agent's tool calls.
4+
5+use g1t_contracts::time::rfc3339;
6+use g1t_contracts::work::*;
7+use g1t_contracts::{FailureCode, Outcome, PrincipalKind, new_id};
8+use g1t_kit::now_ms;
9+use serde::Deserialize;
10+use worker::Result;
11+
12+use crate::Work;
13+
14+/// The longest message an agent is sent.
15+const MAX_MESSAGE_CHARS: usize = 4000;
16+
17+#[derive(Deserialize)]
18+struct MessageRow {
19+ id: String,
20+ author_name: String,
21+ body: String,
22+ created_at: String,
23+ delivered_at: Option<String>,
24+}
25+
26+impl From<MessageRow> for AgentMessage {
27+ fn from(row: MessageRow) -> Self {
28+ AgentMessage {
29+ id: row.id,
30+ author: row.author_name,
31+ body: row.body,
32+ created_at: row.created_at,
33+ delivered_at: row.delivered_at,
34+ }
35+ }
36+}
37+
38+impl Work {
39+ /// Every message sent to the agent on a pull request, oldest first.
40+ pub(crate) async fn messages(&self, pull_id: &str) -> Result<Vec<AgentMessage>> {
41+ Ok(self
42+ .db
43+ .prepare("SELECT * FROM agent_messages WHERE pull_id = ? ORDER BY created_at, id")
44+ .bind(&[pull_id.into()])?
45+ .all()
46+ .await?
47+ .results::<MessageRow>()?
48+ .into_iter()
49+ .map(AgentMessage::from)
50+ .collect())
51+ }
52+
53+ pub(crate) async fn message_agent(&self, a: MessageAgentArgs) -> Result<Outcome<AgentMessage>> {
54+ let viewer = Some(a.actor.clone());
55+ let (repo, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
56+ Outcome::Ok(found) => found,
57+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
58+ };
59+ if !a.actor.verified
60+ || (pull.author.id != a.actor.id && !a.actor.is_member(&repo.namespace))
61+ {
62+ return Ok(Outcome::fail(
63+ FailureCode::Forbidden,
64+ "Only the pull request's author and members of the workspace can message its agent.",
65+ ));
66+ }
67+ if !pull.status.is_active() {
68+ return Ok(Outcome::fail(
69+ FailureCode::Conflict,
70+ "This pull request is no longer being worked on.",
71+ ));
72+ }
73+ let body = a.body.trim();
74+ if body.is_empty() {
75+ return Ok(Outcome::fail(FailureCode::Invalid, "Write a message."));
76+ }
77+ let body: String = body.chars().take(MAX_MESSAGE_CHARS).collect();
78+ let now = now_ms();
79+ let message = AgentMessage {
80+ id: new_id("msg", now),
81+ author: a.actor.username.clone(),
82+ body,
83+ created_at: rfc3339(now),
84+ delivered_at: None,
85+ };
86+ self.db
87+ .prepare(
88+ "INSERT INTO agent_messages (id, pull_id, author_id, author_name, body, created_at)
89+ VALUES (?, ?, ?, ?, ?, ?)",
90+ )
91+ .bind(&[
92+ message.id.as_str().into(),
93+ pull.id.as_str().into(),
94+ a.actor.id.as_str().into(),
95+ message.author.as_str().into(),
96+ message.body.as_str().into(),
97+ message.created_at.as_str().into(),
98+ ])?
99+ .run()
100+ .await?;
101+ self.note(
102+ &repo.id,
103+ pull.number,
104+ (a.actor.id.as_str(), a.actor.username.as_str()),
105+ "sent the agent a message",
106+ )
107+ .await?;
108+ Ok(Outcome::Ok(message))
109+ }
110+
111+ /// The undelivered messages, marked delivered and recorded in the
112+ /// session, for the agent's sandbox.
113+ pub(crate) async fn take_messages(&self, a: TakeMessagesArgs) -> Result<Outcome<Vec<AgentMessage>>> {
114+ if a.actor.kind != PrincipalKind::Agent {
115+ return Ok(Outcome::fail(
116+ FailureCode::Forbidden,
117+ "Only g1t's agents take messages.",
118+ ));
119+ }
120+ let viewer = Some(a.actor.clone());
121+ let (_, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
122+ Outcome::Ok(found) => found,
123+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
124+ };
125+ let now = rfc3339(now_ms());
126+ let taken: Vec<AgentMessage> = self
127+ .db
128+ .prepare(
129+ "UPDATE agent_messages SET delivered_at = ?
130+ WHERE pull_id = ? AND delivered_at IS NULL
131+ RETURNING *",
132+ )
133+ .bind(&[now.as_str().into(), pull.id.as_str().into()])?
134+ .all()
135+ .await?
136+ .results::<MessageRow>()?
137+ .into_iter()
138+ .map(AgentMessage::from)
139+ .collect();
140+ if !taken.is_empty() {
141+ let entries: Vec<NewSessionEntry> = taken
142+ .iter()
143+ .map(|message| NewSessionEntry {
144+ kind: SessionEntryKind::Prompt,
145+ text: format!("Message from {}: {}", message.author, message.body),
146+ tool: None,
147+ commit: None,
148+ })
149+ .collect();
150+ self.append_entries(&pull, &entries).await?;
151+ }
152+ Ok(Outcome::Ok(taken))
153+ }
154+}