Skip to content
133 linesCodeBlameRaw
1/**
2 * The skill library's RPC methods (@g1t/contracts skill-library.ts
3 * `skillLibraryClient`), and what the library reads from identity, repos
4 * and chat for them. The rules are in skill-library.ts.
5 */
6import { type RepoPath, type Result, type User, chatClient, fail, identityClient, ok, reposClient } from "@g1t/contracts";
7
8import { metered } from "./meter.ts";
9import { type SessionEnv, sessionRow } from "./sessions.ts";
10import { saveDraft, DRAFT_SYSTEM, transcriptText } from "./skill-draft.ts";
11import { Library, type LibraryPorts, onPush } from "./skill-library.ts";
12import type { Row } from "./store.ts";
13import { runTurn } from "./turn.ts";
14import type { ViewContext } from "./views.ts";
15
16/** Reads repositories without a viewer: only after the viewer's access was checked, or for a linked repository. */
17function repoFiles(env: { REPOS: SessionEnv["REPOS"] }): Pick<LibraryPorts, "listFiles" | "blobs"> {
18 const repos = reposClient(env.REPOS);
19 return {
20 listFiles: (repoId, ref) => repos.listFiles(repoId, ref, 10_000),
21 blobs: async (repoId, hashes) => (await repos.rawBlobs(repoId, hashes, 1024 * 1024)).map((blob) => ({ hash: blob.hash, data: blob.data })),
22 };
23}
24
25export function libraryFor(ctx: ViewContext, audit: (action: string, name: string, message: string) => void): Library {
26 const env = ctx.env;
27 const identity = identityClient(env.IDENTITY);
28 const ports: LibraryPorts = {
29 teams: async () => {
30 const listed = await identity.listTeams(ctx.viewer, ctx.slug).catch(() => null);
31 return listed?.ok ? listed.value.map((t) => ({ slug: t.slug, name: t.name, can_manage: t.can_manage })) : null;
32 },
33 repo: async (full) => {
34 const [namespace, name] = full.split("/") as [string, string];
35 const found = await reposClient(env.REPOS)
36 .get({ namespace, name } as RepoPath, ctx.viewer)
37 .catch(() => null);
38 return found?.ok ? { id: found.value.id, full: `${found.value.namespace}/${found.value.name}`, default_branch: found.value.defaultBranch } : null;
39 },
40 ...repoFiles(env),
41 agentTeams: async (agentId) => (await identity.agentTeams(ctx.slug, agentId)).map((t) => ({ slug: t.slug, name: t.name })),
42 teamAgentIndex: async () => (await identity.teamAgentIndex(ctx.slug)).map((t) => ({ slug: t.slug, agent_ids: t.agent_ids })),
43 audit,
44 };
45 return new Library({ db: ctx.db, workspaceId: ctx.workspaceId, slug: ctx.slug, viewer: { id: ctx.viewer.id, username: ctx.viewer.username, kind: ctx.viewer.kind }, owner: ctx.owner, ports });
46}
47
48/** Save as skill: the session's agent drafts it from the transcript, on its own budget. */
49async function draftSkill(ctx: ViewContext, library: Library, id: unknown): Promise<Result<unknown>> {
50 if ((ctx.viewer.kind ?? "user") !== "user") return fail("forbidden", "People save sessions as skills.");
51 const row = await sessionRow(ctx.db, String(id ?? ""));
52 if (!row || row.workspace_id !== ctx.workspaceId) return fail("not_found", "There is no such session.");
53 const audience = await chatClient(ctx.env.CHAT)
54 .audience(ctx.slug, row.channel_id)
55 .catch(() => null);
56 const visible = !!audience?.ok && (audience.value.kind === "public" || audience.value.member_user_ids.includes(ctx.viewer.id));
57 if (!visible) return fail("not_found", "There is no such session.");
58 if (row.status !== "done") return fail("invalid", "Save a session as a skill once it is done.");
59 const agent = await ctx.db.prepare("SELECT * FROM agents WHERE id = ?").bind(row.agent_id).first<Row>();
60 if (!agent) return fail("not_found", "The session's agent is gone.");
61 const events = await ctx.db
62 .prepare("SELECT kind, by_name, body, tool FROM agent_session_events WHERE session_id = ? ORDER BY seq LIMIT 600")
63 .bind(row.id)
64 .all<{ kind: string; by_name: string | null; body: string; tool: string | null }>();
65 const transcript = transcriptText({ title: row.title, goal: row.goal, result: row.summary }, events.results);
66 return saveDraft(library, { id: row.id, title: row.title, agent_handle: agent.handle }, async () => {
67 const done = await metered(
68 ctx.env,
69 { row: agent, payer: agent, slug: ctx.slug, task: "session", start: "small", askerName: ctx.viewer.username, person: ctx.viewer.username },
70 async (model) => {
71 const answer = await runTurn(model.send, {
72 model: model.model.model,
73 system: DRAFT_SYSTEM,
74 messages: [{ role: "user", content: transcript }],
75 tools: null,
76 price: model.ownModel ? null : model.model.price,
77 maxRounds: 1,
78 maxOutput: 4096,
79 });
80 return { ...answer, cost: model.ownModel ? 0 : answer.cost };
81 },
82 ).catch((error: unknown) => ({ ok: false as const, reason: "error", message: error instanceof Error ? error.message : String(error) }));
83 if (!done.ok) return fail(done.reason === "error" ? "unavailable" : "limit", done.reason === "error" ? "The draft couldn't be written just now. Try again in a moment." : done.message);
84 return ok(done.value.text);
85 });
86}
87
88/** The library's methods, or null for one that isn't the library's. */
89export async function skillRpc(
90 method: string,
91 args: any,
92 view: <T>(a: { workspace: string; viewer: User | null }, run: (ctx: ViewContext) => Promise<Result<T>>) => Promise<Result<T>>,
93 audit: (viewer: User, workspace: string, action: string, name: string, message: string) => void,
94): Promise<Result<unknown> | null> {
95 const library = (ctx: ViewContext) => libraryFor(ctx, (action, name, message) => audit(ctx.viewer, ctx.slug, action, name, message));
96 switch (method) {
97 case "skill_library":
98 return view(args, (ctx) => library(ctx).library());
99 case "skill":
100 return view(args, (ctx) => library(ctx).detail(args.name, args.version));
101 case "save_skill":
102 return view(args, (ctx) => library(ctx).save(args.name, args.input));
103 case "import_skill":
104 return view(args, (ctx) => library(ctx).import(args.source, args.replace === true));
105 case "attach_skill":
106 return view(args, (ctx) => library(ctx).attach(args.name, args.scope, args.target));
107 case "detach_skill":
108 return view(args, (ctx) => library(ctx).detach(args.name, args.attachment));
109 case "pin_skill":
110 return view(args, (ctx) => library(ctx).pin(args.name, args.attachment, args.version ?? null));
111 case "delete_skill":
112 return view(args, (ctx) => library(ctx).remove(args.name));
113 case "agent_skills":
114 return view(args, (ctx) => library(ctx).agentSkills(args.handle));
115 case "draft_skill":
116 return view(args, (ctx) => draftSkill(ctx, library(ctx), args.session));
117 case "set_skill_mirror":
118 return view(args, (ctx) => library(ctx).setMirror(args.repo ?? null));
119 case "sync_skill_mirror":
120 return view(args, (ctx) => library(ctx).sync());
121 default:
122 return null;
123 }
124}
125
126/** After pushes: libraries that follow a pushed repository's default branch read it again. */
127export async function skillPushes(env: { DB: D1Database; REPOS: SessionEnv["REPOS"] }, pushed: { repoId: string }[]): Promise<void> {
128 const repos = [...new Set(pushed.map((p) => p.repoId))];
129 for (const repoId of repos) {
130 await onPush(env.DB, repoFiles(env), repoId).catch((error: unknown) => console.error("agents: skills not read after a push", repoId, String(error)));
131 }
132}
133