Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/good-job/implementation.ts

Raw
import { spawn } from "node:child_process";
import { existsSync } from "node:fs";
import { basename, dirname, join } from "node:path";
import type {
	AssistantMessage,
	Model,
	ThinkingLevel,
} from "@earendil-works/pi-ai";
import {
	buildSessionContext,
	createAgentSession,
	createEventBus,
	DefaultResourceLoader,
	defineTool,
	type ExtensionAPI,
	type ExtensionContext,
	type FileEntry,
	getAgentDir,
	ModelRuntime,
	type SessionEntry,
	SessionManager,
	SettingsManager,
} from "@earendil-works/pi-coding-agent";
import {
	analysisPrompt,
	boundedError,
	type FeedbackKind,
	feedbackId,
	type GoodJobJob,
	type GoodJobRecord,
	JOB_VERSION,
	type Learning,
	LearningSubmissionSchema,
	lastAgentModel,
	type ModelRef,
	parseLearning,
	RECORD_VERSION,
} from "./core.js";
import { dbg } from "./src/debug.ts";
import { findExecutable } from "./src/pi-ext-executable.ts";
import {
	resolveSetting,
	type SettingDeclaration,
} from "./src/pi-ext-settings.ts";
import {
	claimJob,
	commitFailure,
	commitRecord,
	configuredDataDirectory,
	counts,
	enqueueJob,
	ensureLayout,
	type GoodJobLayout,
	goodJobDataDirectory,
	layout,
	loadRecords,
	migrateStore,
	pendingJobIds,
	recoverProcessing,
} from "./storage.js";
import { showLearningDetail, showLearningList } from "./ui.js";

const STATUS_KEY = "good-job";
const SPINNER = ["◐", "◓", "◑", "◒"] as const;
const MAX_TOKENS = 1_200;

type BackgroundContext = {
	modelRegistry: ExtensionContext["modelRegistry"];
	ui?: ExtensionContext["ui"];
};

type Opener = { command: string; args: string[] };
type GoodJobSettings = {
	directory?: string;
	model?: ModelRef;
	thinkingLevel?: ThinkingLevel;
};

const THINKING_LEVELS = new Set<ThinkingLevel>([
	"minimal",
	"low",
	"medium",
	"high",
	"xhigh",
	"max",
]);

export interface GoodJobRuntime {
	start(ctx: ExtensionContext): Promise<void>;
	stop(ctx: ExtensionContext): Promise<void>;
	capture(
		kind: FeedbackKind,
		feedback: string,
		ctx: ExtensionContext,
	): Promise<void>;
	runCommand(args: string, ctx: ExtensionContext): Promise<void>;
}

function modelRef(ctx: ExtensionContext): ModelRef | null {
	return ctx.model ? { provider: ctx.model.provider, id: ctx.model.id } : null;
}

function backgroundContext(ctx: ExtensionContext): BackgroundContext {
	return {
		modelRegistry: ctx.modelRegistry,
		ui: ctx.hasUI ? ctx.ui : undefined,
	};
}

export function parseModelRef(value: string): ModelRef | undefined {
	const trimmed = value.trim();
	const separator = trimmed.indexOf("/");
	if (separator <= 0 || separator === trimmed.length - 1) return undefined;
	return {
		provider: trimmed.slice(0, separator),
		id: trimmed.slice(separator + 1),
	};
}

// Precedence per key: trusted project, user, unset.
const DIRECTORY_SETTING: SettingDeclaration<string | undefined> = {
	key: "good-job.directory",
	parse(raw) {
		if (typeof raw !== "string" || !raw.trim())
			throw new Error("must be a non-empty path");
		return raw.trim();
	},
	default: undefined,
};

const MODEL_SETTING: SettingDeclaration<ModelRef | undefined> = {
	key: "good-job.model",
	parse(raw) {
		const model = typeof raw === "string" ? parseModelRef(raw) : undefined;
		if (!model) throw new Error("must be provider/model");
		return model;
	},
	default: undefined,
};

const THINKING_LEVEL_SETTING: SettingDeclaration<ThinkingLevel | undefined> = {
	key: "good-job.thinkingLevel",
	parse(raw) {
		if (typeof raw === "string" && THINKING_LEVELS.has(raw as ThinkingLevel))
			return raw as ThinkingLevel;
		throw new Error("must be minimal, low, medium, high, xhigh, or max");
	},
	default: undefined,
};

function settingsFrom(
	pi: ExtensionAPI,
	ctx: ExtensionContext,
): GoodJobSettings {
	const read = <T>(declaration: SettingDeclaration<T>): T => {
		const setting = resolveSetting(pi, ctx, declaration);
		if (!setting.ok) throw new Error(setting.error);
		return setting.value;
	};
	return {
		directory: read(DIRECTORY_SETTING),
		model: read(MODEL_SETTING),
		thinkingLevel: read(THINKING_LEVEL_SETTING),
	};
}

function sourceFor(
	ctx: ExtensionContext,
	messages: readonly unknown[],
): GoodJobJob["source"] {
	return {
		type: "session",
		cwd: ctx.cwd,
		sessionPath: ctx.sessionManager.getSessionFile() ?? null,
		sessionId: ctx.sessionManager.getSessionId(),
		sessionName: ctx.sessionManager.getSessionName() ?? null,
		leafId: ctx.sessionManager.getLeafId(),
		agentModel: lastAgentModel(messages),
	};
}

function copyRegisteredProviders(
	registry: ExtensionContext["modelRegistry"],
	runtime: ModelRuntime,
): void {
	for (const providerId of registry.getRegisteredProviderIds()) {
		const native = registry.getRegisteredNativeProvider(providerId);
		if (native) {
			runtime.registerNativeProvider(native);
			continue;
		}
		const config = registry.getRegisteredProviderConfig(providerId);
		if (config) runtime.registerProvider(providerId, config);
	}
}

export function openerFor(
	platform = process.platform,
	env: NodeJS.ProcessEnv = process.env,
): Opener | undefined {
	if (platform === "win32") {
		const root = env.SystemRoot || env.WINDIR;
		const executable = root ? join(root, "explorer.exe") : undefined;
		return executable && existsSync(executable)
			? { command: executable, args: [] }
			: undefined;
	}
	if (platform === "darwin")
		return existsSync("/usr/bin/open")
			? { command: "/usr/bin/open", args: [] }
			: undefined;
	const executable = findExecutable(["xdg-open"], {
		path: env.PATH ?? "",
		platform,
	});
	return executable ? { command: executable, args: [] } : undefined;
}

export function openDirectory(
	path: string,
	opener = openerFor(),
): Promise<void> {
	if (!opener)
		return Promise.reject(new Error("no platform directory opener found"));
	return new Promise((resolve, reject) => {
		const child = spawn(opener.command, [...opener.args, path], {
			detached: true,
			stdio: "ignore",
			windowsHide: true,
		});
		child.once("error", reject);
		child.once("spawn", () => {
			child.unref();
			resolve();
		});
	});
}

function learningPrompt(recordsPath: string): string {
	return `Review every structured GJ learning record under ${JSON.stringify(recordsPath)}.
Read the files yourself; the path is the only learning data supplied in this prompt.
Analyze praise and problem records together.
Reconcile conflicting signals and distinguish recurring evidence from isolated incidents.
Propose concrete changes to Pi configuration, instructions, skills, extensions, or workflows.
Do not edit anything; proposals only.
For every claim and proposal, cite supporting records as \`slug (id)\` using both fields exactly.`;
}

class Runtime implements GoodJobRuntime {
	private generation = 0;
	private activeController: AbortController | undefined;
	private pumpPromise: Promise<void> | undefined;
	private pumpRequested = false;
	private readonly snapshots = new Map<string, SessionEntry[]>();
	private spinnerTimer: ReturnType<typeof setInterval> | undefined;
	private spinnerStartedAt = 0;
	private spinnerCount = 0;
	private openerChecked = false;
	private settings: GoodJobSettings = {};
	private settingsKey: string | undefined;
	private paths: GoodJobLayout;
	private readonly opener: Opener | undefined;
	private readonly defaultRoot = goodJobDataDirectory();

	constructor(
		private readonly pi: ExtensionAPI,
		private readonly fixedRoot?: string,
	) {
		this.paths = layout(fixedRoot ?? this.defaultRoot);
		this.opener = openerFor();
	}

	async start(ctx: ExtensionContext): Promise<void> {
		const background = backgroundContext(ctx);
		const previousPump = this.pumpPromise;
		const generation = ++this.generation;
		this.activeController?.abort();
		if (previousPump) await previousPump;
		if (generation !== this.generation) return;
		this.pumpRequested = false;
		this.readSettings(ctx);
		const targetRoot = this.paths.root;
		const targetMigration = await migrateStore(targetRoot, targetRoot);
		const defaultMigration =
			targetRoot === this.defaultRoot
				? { jobs: 0, records: 0, failures: 0, active: 0 }
				: await migrateStore(this.defaultRoot, targetRoot);
		const migrated = {
			jobs: targetMigration.jobs + defaultMigration.jobs,
			records: targetMigration.records + defaultMigration.records,
			failures: targetMigration.failures + defaultMigration.failures,
			active: targetMigration.active + defaultMigration.active,
		};
		if (migrated.jobs || migrated.records || migrated.failures)
			background.ui?.notify(
				`gj migrated · ${migrated.records} records · ${migrated.jobs} queued · ${migrated.failures} failed`,
				"info",
			);
		if (migrated.active)
			background.ui?.notify(
				`gj left ${migrated.active} active processing claim${migrated.active === 1 ? "" : "s"} in the old store`,
				"warning",
			);
		if (!this.openerChecked) {
			this.openerChecked = true;
			if (!this.opener)
				background.ui?.notify(
					"good-job: directory opening unavailable; no platform opener found",
					"info",
				);
		}
		if (!existsSync(this.paths.root)) return;
		await this.recoverAndPump(background, generation);
	}

	async stop(ctx: ExtensionContext): Promise<void> {
		const ui = ctx.hasUI ? ctx.ui : undefined;
		this.generation++;
		this.pumpRequested = false;
		this.activeController?.abort();
		await this.pumpPromise;
		this.activeController = undefined;
		this.stopSpinner(ui);
	}

	async capture(
		kind: FeedbackKind,
		feedback: string,
		ctx: ExtensionContext,
	): Promise<void> {
		const background = backgroundContext(ctx);
		this.readSettings(ctx);
		const sessionContext = buildSessionContext(
			ctx.sessionManager.getEntries(),
			ctx.sessionManager.getLeafId(),
		);
		const job: GoodJobJob = {
			version: JOB_VERSION,
			id: feedbackId(kind),
			createdAt: new Date().toISOString(),
			kind,
			feedback: feedback.trim(),
			source: sourceFor(ctx, sessionContext.messages),
			analysisModel: this.analysisModel(ctx),
		};
		await ensureLayout(this.paths);
		await enqueueJob(job, this.paths);
		dbg?.("feedback.capture", { kind });
		this.snapshots.set(
			job.id,
			structuredClone(ctx.sessionManager.getEntries()),
		);
		background.ui?.notify(`${job.id} queued`, "info");
		void this.pump(background, this.generation, job.id);
	}

	async runCommand(args: string, ctx: ExtensionContext): Promise<void> {
		this.readSettings(ctx);
		const [command = "list"] = args.trim().split(/\s+/, 1);
		if (command === "list" || command === "") {
			await this.list(ctx);
			return;
		}
		if (command === "status") {
			await this.status(ctx);
			return;
		}
		if (command === "open") {
			await this.open(ctx);
			return;
		}
		if (command === "learn") {
			await this.learn(ctx);
			return;
		}
		ctx.ui.notify("Usage: /gj [list|status|open|learn]", "warning");
	}

	private readSettings(ctx: ExtensionContext): void {
		const key = `${ctx.isProjectTrusted()}:${ctx.cwd}`;
		if (this.settingsKey === key) return;
		this.settings = settingsFrom(this.pi, ctx);
		this.paths = layout(
			this.fixedRoot ??
				configuredDataDirectory(this.settings.directory, ctx.cwd),
		);
		this.settingsKey = key;
	}

	private analysisModel(ctx: ExtensionContext): ModelRef | null {
		return this.settings.model ?? modelRef(ctx);
	}

	private async recoverAndPump(
		ctx: BackgroundContext,
		generation: number,
	): Promise<void> {
		await ensureLayout(this.paths);
		await recoverProcessing(this.paths);
		if (generation === this.generation) void this.pump(ctx, generation);
	}

	private pump(
		ctx: BackgroundContext,
		generation: number,
		preferredId?: string,
	): Promise<void> {
		if (this.pumpPromise) {
			this.pumpRequested = true;
			return this.pumpPromise;
		}
		this.pumpRequested = false;
		this.pumpPromise = this.pumpJobs(ctx, generation, preferredId)
			.catch((error) => {
				if (generation === this.generation)
					ctx.ui?.notify(
						`good-job queue failed: ${boundedError(error)}`,
						"error",
					);
			})
			.finally(() => {
				this.pumpPromise = undefined;
				if (this.pumpRequested && generation === this.generation)
					void this.pump(ctx, generation);
			});
		return this.pumpPromise;
	}

	private async pumpJobs(
		ctx: BackgroundContext,
		generation: number,
		preferredId?: string,
	): Promise<void> {
		let ids = await pendingJobIds(this.paths);
		if (preferredId && ids.includes(preferredId))
			ids = [preferredId, ...ids.filter((id) => id !== preferredId)];
		this.startSpinner(ctx, ids.length);
		try {
			for (const [index, id] of ids.entries()) {
				if (generation !== this.generation) return;
				this.spinnerCount = ids.length - index;
				const claimed = await claimJob(id, this.paths);
				if (!claimed) continue;
				await this.process(claimed.job, claimed.path, ctx, generation);
			}
		} finally {
			this.stopSpinner(ctx.ui);
		}
	}

	private async process(
		job: GoodJobJob,
		claimPath: string,
		ctx: BackgroundContext,
		generation: number,
	): Promise<void> {
		const controller = new AbortController();
		this.activeController = controller;
		try {
			if (!job.analysisModel)
				throw new Error("no analysis model was available");
			const model = ctx.modelRegistry.find(
				job.analysisModel.provider,
				job.analysisModel.id,
			);
			if (!model)
				throw new Error(
					`analysis model ${job.analysisModel.provider}/${job.analysisModel.id} is unavailable`,
				);
			const { learning, response } = await this.runChildAnalysis(
				job,
				model,
				ctx,
				controller.signal,
			);
			if (controller.signal.aborted || generation !== this.generation) {
				dbg?.("learning.outcome", { status: "aborted" });
				return;
			}
			const record: GoodJobRecord = {
				version: RECORD_VERSION,
				id: job.id,
				createdAt: job.createdAt,
				completedAt: new Date().toISOString(),
				kind: job.kind,
				feedback: job.feedback,
				source: job.source,
				analysisModel: job.analysisModel,
				learning,
				usage: {
					input: response.usage.input,
					output: response.usage.output,
					cacheRead: response.usage.cacheRead,
					cacheWrite: response.usage.cacheWrite,
					totalTokens: response.usage.totalTokens,
					cost: response.usage.cost.total,
				},
			};
			await commitRecord(record, claimPath, this.paths);
			dbg?.("learning.outcome", { status: "committed" });
			this.snapshots.delete(job.id);
			ctx.ui?.notify(`${record.id} learned · ${record.learning.slug}`, "info");
		} catch (error) {
			if (controller.signal.aborted || generation !== this.generation) {
				dbg?.("learning.outcome", { status: "aborted" });
				return;
			}
			dbg?.("learning.outcome", { status: "failed" });
			let message = `good-job failed: ${boundedError(error)}`;
			try {
				await commitFailure(job, error, claimPath, this.paths);
				this.snapshots.delete(job.id);
			} catch (storageError) {
				message += `; could not persist failure: ${boundedError(storageError)}`;
			}
			ctx.ui?.notify(message, "error");
		} finally {
			if (this.activeController === controller)
				this.activeController = undefined;
		}
	}

	private async runChildAnalysis(
		job: GoodJobJob,
		model: Model<any>,
		ctx: BackgroundContext,
		signal: AbortSignal,
	): Promise<{ learning: Learning; response: AssistantMessage }> {
		const sourcePath = job.source.sessionPath;
		let manager: SessionManager;
		if (sourcePath) {
			manager = SessionManager.forkFrom(
				sourcePath,
				job.source.cwd,
				dirname(sourcePath),
				{ parentSession: sourcePath },
			);
		} else {
			const snapshot = this.snapshots.get(job.id);
			if (!snapshot)
				throw new Error(
					"source session was ephemeral and its snapshot is unavailable",
				);
			const header: FileEntry = {
				type: "session",
				version: 3,
				id: job.source.sessionId,
				timestamp: job.createdAt,
				cwd: job.source.cwd,
			};
			manager = SessionManager.inMemory(job.source.cwd, undefined, [
				header,
				...snapshot,
			]);
		}
		if (job.source.leafId) {
			if (!manager.getEntry(job.source.leafId))
				throw new Error(`source session leaf ${job.source.leafId} is missing`);
			manager.branch(job.source.leafId);
		}
		manager.appendSessionInfo(`good-job · ${job.id}`);

		const agentDir = getAgentDir();
		const childRuntime = await ModelRuntime.create({
			authPath: join(agentDir, "auth.json"),
			modelsPath: join(agentDir, "models.json"),
			refreshOnCreate: false,
			signal,
		});
		copyRegisteredProviders(ctx.modelRegistry, childRuntime);

		const settingsManager = SettingsManager.create(job.source.cwd, agentDir);
		const resourceLoader = new DefaultResourceLoader({
			cwd: job.source.cwd,
			agentDir,
			settingsManager,
			eventBus: createEventBus(),
			noExtensions: true,
			noSkills: true,
			noPromptTemplates: true,
			noThemes: true,
			noContextFiles: true,
		});
		await resourceLoader.reload();
		signal.throwIfAborted();

		const childModel = {
			...model,
			maxTokens:
				model.maxTokens > 0
					? Math.min(model.maxTokens, MAX_TOKENS)
					: MAX_TOKENS,
		};
		let learning: Learning | undefined;
		let duplicateSubmission = false;
		const submitLearningTool = defineTool<
			typeof LearningSubmissionSchema,
			Learning
		>({
			name: "submit_learning",
			label: "Submit Learning",
			description: "Submit the final structured learning and finish analysis.",
			parameters: LearningSubmissionSchema,
			constrainedSampling: { type: "json_schema", strict: "prefer" },
			async execute(_toolCallId, params) {
				if (learning) {
					duplicateSubmission = true;
					throw new Error("submit_learning may be called exactly once");
				}
				const parsed = parseLearning(JSON.stringify(params));
				learning = parsed;
				return {
					content: [{ type: "text", text: "Learning submitted." }],
					details: parsed,
					terminate: true,
				};
			},
		});
		let child:
			| Awaited<ReturnType<typeof createAgentSession>>["session"]
			| undefined;
		try {
			const created = await createAgentSession({
				cwd: job.source.cwd,
				agentDir,
				model: childModel,
				thinkingLevel: this.settings.thinkingLevel,
				modelRuntime: childRuntime,
				resourceLoader,
				settingsManager,
				sessionManager: manager,
				tools: [submitLearningTool.name],
				customTools: [submitLearningTool],
			});
			child = created.session;
			const abortChild = () => void child?.abort();
			signal.addEventListener("abort", abortChild, { once: true });
			try {
				await child.prompt(analysisPrompt(job.kind, job.feedback), {
					expandPromptTemplates: false,
					source: "extension",
				});
			} finally {
				signal.removeEventListener("abort", abortChild);
			}
			signal.throwIfAborted();
			const finalAssistant = [...child.messages]
				.reverse()
				.find(
					(message): message is AssistantMessage =>
						message.role === "assistant",
				);
			if (!finalAssistant)
				throw new Error(
					"child analysis finished without an assistant response",
				);
			if (duplicateSubmission)
				throw new Error("analysis called submit_learning more than once");
			if (!learning)
				throw new Error(
					finalAssistant.errorMessage ||
						"analysis finished without calling submit_learning",
				);
			return { learning, response: finalAssistant };
		} finally {
			if (child) {
				if (!child.isIdle) await child.abort();
				child.dispose();
			}
		}
	}

	private startSpinner(ctx: BackgroundContext, count: number): void {
		const { ui } = ctx;
		this.stopSpinner(ui);
		if (!ui || count === 0) return;
		this.spinnerCount = count;
		this.spinnerStartedAt = Date.now();
		const render = () => {
			const frame =
				SPINNER[
					Math.floor((Date.now() - this.spinnerStartedAt) / 120) %
						SPINNER.length
				] ?? SPINNER[0];
			ui.setStatus(
				STATUS_KEY,
				ui.theme.fg(
					"accent",
					`${frame} good-job learning${this.spinnerCount > 1 ? ` (${this.spinnerCount})` : ""}`,
				),
			);
		};
		render();
		this.spinnerTimer = setInterval(render, 120);
		this.spinnerTimer.unref?.();
	}

	private stopSpinner(ui?: ExtensionContext["ui"]): void {
		if (this.spinnerTimer) clearInterval(this.spinnerTimer);
		this.spinnerTimer = undefined;
		ui?.setStatus(STATUS_KEY, undefined);
	}

	private async list(ctx: ExtensionContext): Promise<void> {
		const loaded = await loadRecords(this.paths);
		if (loaded.errors.length)
			ctx.ui.notify(
				`${loaded.errors.length} unreadable good-job record${loaded.errors.length === 1 ? "" : "s"}`,
				"warning",
			);
		if (loaded.records.length === 0) {
			ctx.ui.notify("No good-job learnings yet", "info");
			return;
		}
		const selected = await showLearningList(
			ctx,
			loaded.records,
			await counts(this.paths),
		);
		if (!selected) return;
		const record = loaded.records.find(
			(candidate) => candidate.id === selected,
		);
		if (record) await showLearningDetail(ctx, record);
	}

	private async status(ctx: ExtensionContext): Promise<void> {
		const state = await counts(this.paths);
		ctx.ui.notify(
			[
				`${state.records} learnings`,
				`${state.pending} pending`,
				`${state.processing} processing`,
				`${state.failed} failed`,
				`analysis model: ${this.settings.model ? `${this.settings.model.provider}/${this.settings.model.id}` : "current"}`,
				`thinking: ${this.settings.thinkingLevel ?? "provider default"}`,
				this.paths.root,
			].join(" · "),
			state.failed ? "warning" : "info",
		);
	}

	private async learn(ctx: ExtensionContext): Promise<void> {
		const state = await counts(this.paths);
		if (state.records === 0) {
			ctx.ui.notify("No GJ learnings yet", "info");
			return;
		}
		this.pi.sendUserMessage(learningPrompt(this.paths.records), {
			deliverAs: ctx.isIdle() ? undefined : "followUp",
		});
	}

	private async open(ctx: ExtensionContext): Promise<void> {
		await ensureLayout(this.paths);
		try {
			await openDirectory(this.paths.root, this.opener);
		} catch (error) {
			ctx.ui.notify(`good-job open failed: ${boundedError(error)}`, "error");
			return;
		}
		ctx.ui.notify(`Opened ${basename(this.paths.root)}`, "info");
	}
}

export function createGoodJobRuntime(
	pi: ExtensionAPI,
	root?: string,
): GoodJobRuntime {
	return new Runtime(pi, root);
}