Parent directory

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}