herdr-agent-state.ts
6694 bytes
1// installed by herdr
2// managed by herdr; reinstalling or updating the integration overwrites this file.
3// add custom hooks/plugins beside this file instead of editing it.
4// HERDR_INTEGRATION_ID=pi
5// HERDR_INTEGRATION_VERSION=9
6// @ts-nocheck
7
8import net from "node:net";
9import path from "node:path";
10
11const HERDR_ENV = process.env.HERDR_ENV;
12const socketPath = process.env.HERDR_SOCKET_PATH;
13const socketEndpoint =
14 process.platform === "win32" && socketPath ? `\\\\.\\pipe\\${socketPath}` : socketPath;
15const paneId = process.env.HERDR_PANE_ID;
16const source = "herdr:pi";
17
18function enabled() {
19 return HERDR_ENV === "1" && !!socketPath && !!paneId;
20}
21
22function sendRequestAttempt(request: unknown, timeoutMs: number): Promise<boolean> {
23 if (!enabled()) {
24 return Promise.resolve(true);
25 }
26
27 return new Promise((resolve) => {
28 let done = false;
29 let timeout: ReturnType<typeof setTimeout> | undefined;
30 const finish = (delivered: boolean) => {
31 if (done) return;
32 done = true;
33 if (timeout) {
34 clearTimeout(timeout);
35 }
36 socket.destroy();
37 resolve(delivered);
38 };
39
40 const socket = net.createConnection(socketEndpoint!);
41 socket.on("error", () => finish(false));
42 socket.on("connect", () => socket.write(`${JSON.stringify(request)}\n`));
43 socket.on("data", () => finish(true));
44 socket.on("end", () => finish(false));
45 timeout = setTimeout(() => finish(false), timeoutMs);
46 timeout.unref?.();
47 });
48}
49
50async function sendRequest(request: unknown): Promise<void> {
51 if (await sendRequestAttempt(request, 500)) {
52 return;
53 }
54 await sendRequestAttempt(request, 1500);
55}
56
57type AgentState = "working" | "blocked" | "idle";
58
59type QueuedState = {
60 state: AgentState;
61 message?: string;
62 seq: number;
63};
64
65let reportSeq = Date.now() * 1000;
66let currentAgentSessionId: string | undefined;
67let currentAgentSessionPath: string | undefined;
68
69function nextReportSeq(): number {
70 reportSeq += 1;
71 return reportSeq;
72}
73
74function updateSessionRef(ctx: any): void {
75 try {
76 const file = ctx?.sessionManager?.getSessionFile?.();
77 currentAgentSessionPath =
78 typeof file === "string" &&
79 (path.posix.isAbsolute(file) || path.win32.isAbsolute(file))
80 ? file
81 : undefined;
82 } catch {
83 currentAgentSessionPath = undefined;
84 }
85
86 try {
87 const id = ctx?.sessionManager?.getSessionId?.();
88 currentAgentSessionId = typeof id === "string" && id.length > 0 ? id : undefined;
89 } catch {
90 currentAgentSessionId = undefined;
91 }
92}
93
94function withSessionRef(params: Record<string, unknown>): Record<string, unknown> {
95 if (currentAgentSessionPath) {
96 return { ...params, agent_session_path: currentAgentSessionPath };
97 }
98 if (currentAgentSessionId) {
99 return { ...params, agent_session_id: currentAgentSessionId };
100 }
101 return params;
102}
103
104function currentSessionRef(): Record<string, unknown> | undefined {
105 if (currentAgentSessionPath) {
106 return { agent_session_path: currentAgentSessionPath };
107 }
108 if (currentAgentSessionId) {
109 return { agent_session_id: currentAgentSessionId };
110 }
111 return undefined;
112}
113
114function reportSession(sessionStartSource?: string): Promise<void> {
115 const sessionRef = currentSessionRef();
116 if (!sessionRef) {
117 return Promise.resolve();
118 }
119
120 return sendRequest({
121 id: `${source}:session:${Date.now()}:${Math.random().toString(36).slice(2)}`,
122 method: "pane.report_agent_session",
123 params: {
124 pane_id: paneId,
125 source,
126 agent: "pi",
127 seq: nextReportSeq(),
128 session_start_source: sessionStartSource,
129 ...sessionRef,
130 },
131 });
132}
133
134function sendState(state: AgentState, message?: string, seq = nextReportSeq()): Promise<void> {
135 return sendRequest({
136 id: `${source}:${Date.now()}:${Math.random().toString(36).slice(2)}`,
137 method: "pane.report_agent",
138 params: withSessionRef({
139 pane_id: paneId,
140 source,
141 agent: "pi",
142 state,
143 message,
144 seq,
145 }),
146 });
147}
148
149let sendInFlight = false;
150let queuedState: QueuedState | undefined;
151
152function queueState(state: AgentState, message?: string): void {
153 queuedState = { state, message, seq: nextReportSeq() };
154 if (!sendInFlight) {
155 void drainStateQueue();
156 }
157}
158
159async function drainStateQueue(): Promise<void> {
160 if (sendInFlight) {
161 return;
162 }
163
164 sendInFlight = true;
165 try {
166 while (queuedState) {
167 const next = queuedState;
168 queuedState = undefined;
169 await sendState(next.state, next.message, next.seq);
170 }
171 } finally {
172 sendInFlight = false;
173 if (queuedState) {
174 void drainStateQueue();
175 }
176 }
177}
178
179export default function (pi) {
180 if (!enabled()) {
181 return;
182 }
183
184 let agentActive = false;
185 let blockedCount = 0;
186 let blockedMessage: string | undefined;
187 let lastState: AgentState | undefined;
188 let lastMessage: string | undefined;
189 let rootSession = false;
190
191 function desiredState() {
192 if (blockedCount > 0) {
193 return { state: "blocked" as const, message: blockedMessage };
194 }
195 if (agentActive) {
196 return { state: "working" as const, message: undefined };
197 }
198 return { state: "idle" as const, message: undefined };
199 }
200
201 function publishState(force = false) {
202 const next = desiredState();
203 if (!force && next.state === lastState && next.message === lastMessage) {
204 return;
205 }
206 lastState = next.state;
207 lastMessage = next.message;
208 queueState(next.state, next.message);
209 }
210
211 pi.events.on("herdr:blocked", (data) => {
212 if (!rootSession) {
213 return;
214 }
215 if (!data?.active) {
216 blockedCount = Math.max(0, blockedCount - 1);
217 if (blockedCount === 0) {
218 blockedMessage = undefined;
219 }
220 publishState();
221 return;
222 }
223
224 blockedCount += 1;
225 blockedMessage = data.label;
226 publishState();
227 });
228
229 pi.on("session_start", async (event, ctx) => {
230 // TUI only: RPC/JSON/print modes are headless (no PTY herdr can display),
231 // and RPC still reports hasUI=true, so mode is the reliable gate.
232 if (ctx?.mode !== "tui") {
233 return;
234 }
235 rootSession = true;
236 updateSessionRef(ctx);
237 await reportSession(event?.reason);
238 // A reload can replace this extension mid-run without emitting another agent_start.
239 agentActive = ctx?.isIdle?.() === false;
240 publishState(true);
241 });
242
243 pi.on("agent_start", (_event, ctx) => {
244 if (!rootSession) {
245 return;
246 }
247 updateSessionRef(ctx);
248 void reportSession();
249 agentActive = true;
250 publishState();
251 });
252
253 pi.on("agent_settled", (_event, ctx) => {
254 if (!rootSession || ctx?.isIdle?.() !== true) {
255 return;
256 }
257
258 agentActive = false;
259 publishState();
260 });
261}