| 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 | */ |
| 6 | import { type RepoPath, type Result, type User, chatClient, fail, identityClient, ok, reposClient } from "@g1t/contracts"; |
| 7 | |
| 8 | import { metered } from "./meter.ts"; |
| 9 | import { type SessionEnv, sessionRow } from "./sessions.ts"; |
| 10 | import { saveDraft, DRAFT_SYSTEM, transcriptText } from "./skill-draft.ts"; |
| 11 | import { Library, type LibraryPorts, onPush } from "./skill-library.ts"; |
| 12 | import type { Row } from "./store.ts"; |
| 13 | import { runTurn } from "./turn.ts"; |
| 14 | import 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. */ |
| 17 | function 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 | |
| 25 | export 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. */ |
| 49 | async 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. */ |
| 89 | export 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. */ |
| 127 | export 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 | |