mirror of
https://github.com/tale/headplane.git
synced 2026-07-28 00:28:56 +00:00
feat: i guess we're undoing agent work
This commit is contained in:
+106
-43
@@ -1,7 +1,7 @@
|
||||
import { execFile } from "node:child_process";
|
||||
import { type ChildProcess, spawn } from "node:child_process";
|
||||
import { access, constants, mkdir, rm, stat } from "node:fs/promises";
|
||||
import { join } from "node:path";
|
||||
import { promisify } from "node:util";
|
||||
import { createInterface } from "node:readline";
|
||||
|
||||
import { inArray, notInArray } from "drizzle-orm";
|
||||
import { NodeSQLiteDatabase } from "drizzle-orm/node-sqlite";
|
||||
@@ -13,8 +13,6 @@ import { HeadplaneConfig } from "./config/config-schema";
|
||||
import { ephemeralNodes, hostInfo } from "./db/schema";
|
||||
import { RuntimeApiClient } from "./headscale/api/endpoints";
|
||||
|
||||
const execFileAsync = promisify(execFile);
|
||||
|
||||
export interface AgentManager {
|
||||
lookup(nodeKeys: string[]): Promise<Record<string, HostInfo>>;
|
||||
lastSync(): { syncedAt: Date | null; nodeCount: number; error?: string };
|
||||
@@ -26,6 +24,7 @@ export interface AgentManager {
|
||||
interface AgentOutput {
|
||||
self: string;
|
||||
hosts: Record<string, HostInfo>;
|
||||
error?: string;
|
||||
}
|
||||
|
||||
interface SyncState {
|
||||
@@ -33,8 +32,6 @@ interface SyncState {
|
||||
nodeCount: number;
|
||||
selfKey?: string;
|
||||
error?: string;
|
||||
isSyncing: boolean;
|
||||
pendingResync: boolean;
|
||||
}
|
||||
|
||||
async function hasExistingState(workDir: string): Promise<boolean> {
|
||||
@@ -95,10 +92,13 @@ export async function createAgentManager(
|
||||
const state: SyncState = {
|
||||
syncedAt: null,
|
||||
nodeCount: 0,
|
||||
isSyncing: false,
|
||||
pendingResync: false,
|
||||
};
|
||||
|
||||
let proc: ChildProcess | null = null;
|
||||
let responseHandler: ((line: string) => void) | null = null;
|
||||
let disposed = false;
|
||||
let consecutiveErrors = 0;
|
||||
|
||||
async function generateAuthKey(): Promise<string> {
|
||||
const expiration = new Date(Date.now() + 5 * 60_000);
|
||||
const pak = await apiClient.createPreAuthKey(null, false, false, expiration, [
|
||||
@@ -107,7 +107,7 @@ export async function createAgentManager(
|
||||
return pak.key;
|
||||
}
|
||||
|
||||
async function runAgent(authKey: string): Promise<string> {
|
||||
function spawnAgent(authKey: string): ChildProcess {
|
||||
const env: Record<string, string> = {
|
||||
HOME: process.env.HOME ?? "",
|
||||
HEADPLANE_AGENT_WORK_DIR: workDir,
|
||||
@@ -120,53 +120,105 @@ export async function createAgentManager(
|
||||
env.HEADPLANE_AGENT_TS_AUTHKEY = authKey;
|
||||
}
|
||||
|
||||
const { stdout } = await execFileAsync(executablePath, [], {
|
||||
timeout: 60_000,
|
||||
const child = spawn(executablePath, [], {
|
||||
env,
|
||||
stdio: ["pipe", "pipe", "pipe"],
|
||||
});
|
||||
|
||||
return stdout;
|
||||
child.stderr?.on("data", (chunk: Buffer) => {
|
||||
const text = chunk.toString().trim();
|
||||
if (text) {
|
||||
log.debug("agent", "%s", text);
|
||||
}
|
||||
});
|
||||
|
||||
const rl = createInterface({ input: child.stdout! });
|
||||
rl.on("line", (line) => {
|
||||
if (responseHandler) {
|
||||
const handler = responseHandler;
|
||||
responseHandler = null;
|
||||
handler(line);
|
||||
}
|
||||
});
|
||||
|
||||
child.on("exit", (code, signal) => {
|
||||
if (!disposed) {
|
||||
log.warn("agent", "Agent process exited (code=%s, signal=%s)", code, signal);
|
||||
}
|
||||
proc = null;
|
||||
|
||||
// Reject any pending sync request
|
||||
if (responseHandler) {
|
||||
const handler = responseHandler;
|
||||
responseHandler = null;
|
||||
handler("");
|
||||
}
|
||||
});
|
||||
|
||||
proc = child;
|
||||
return child;
|
||||
}
|
||||
|
||||
async function ensureProcess(): Promise<ChildProcess> {
|
||||
if (proc && proc.exitCode === null) {
|
||||
return proc;
|
||||
}
|
||||
|
||||
const stateExists = await hasExistingState(workDir);
|
||||
if (stateExists) {
|
||||
log.debug("agent", "Reusing existing tsnet identity");
|
||||
return spawnAgent("");
|
||||
}
|
||||
|
||||
log.info("agent", "No tsnet state found, generating pre-auth key");
|
||||
return spawnAgent(await generateAuthKey());
|
||||
}
|
||||
|
||||
function sendSync(child: ChildProcess): Promise<string> {
|
||||
return new Promise((resolve) => {
|
||||
responseHandler = resolve;
|
||||
child.stdin?.write("sync\n");
|
||||
});
|
||||
}
|
||||
|
||||
async function requestSync(child: ChildProcess): Promise<AgentOutput> {
|
||||
const line = await sendSync(child);
|
||||
if (!line) {
|
||||
throw new Error("Agent process closed unexpectedly");
|
||||
}
|
||||
return JSON.parse(line) as AgentOutput;
|
||||
}
|
||||
|
||||
let isSyncing = false;
|
||||
let pendingResync = false;
|
||||
|
||||
async function sync() {
|
||||
if (state.isSyncing) {
|
||||
state.pendingResync = true;
|
||||
if (isSyncing) {
|
||||
pendingResync = true;
|
||||
log.debug("agent", "Sync already in progress, queued resync");
|
||||
return;
|
||||
}
|
||||
|
||||
state.isSyncing = true;
|
||||
isSyncing = true;
|
||||
try {
|
||||
const stateExists = await hasExistingState(workDir);
|
||||
const authKey = stateExists ? "" : await generateAuthKey();
|
||||
const child = await ensureProcess();
|
||||
const output = await requestSync(child);
|
||||
|
||||
if (stateExists) {
|
||||
log.debug("agent", "Reusing existing tsnet identity");
|
||||
}
|
||||
if (output.error) {
|
||||
consecutiveErrors++;
|
||||
state.error = output.error;
|
||||
log.error("agent", "Sync error from agent (%d/5): %s", consecutiveErrors, output.error);
|
||||
|
||||
let stdout: string;
|
||||
try {
|
||||
stdout = await runAgent(authKey);
|
||||
} catch (err) {
|
||||
if (!stateExists) {
|
||||
throw err;
|
||||
}
|
||||
|
||||
// Retry once with existing state (e.g. stale lock from a previous run)
|
||||
log.info("agent", "Agent failed with existing state, retrying");
|
||||
try {
|
||||
const retryKey = await generateAuthKey();
|
||||
stdout = await runAgent(retryKey);
|
||||
} catch {
|
||||
// Only clear identity as a last resort
|
||||
log.warn("agent", "Retry failed, clearing state and starting fresh");
|
||||
if (consecutiveErrors >= 5 && proc) {
|
||||
log.warn("agent", "Too many consecutive errors, killing agent and clearing state");
|
||||
proc.kill("SIGTERM");
|
||||
proc = null;
|
||||
await rm(join(workDir, "tailscaled.state"), { force: true });
|
||||
const freshKey = await generateAuthKey();
|
||||
stdout = await runAgent(freshKey);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
const output = JSON.parse(stdout) as AgentOutput;
|
||||
consecutiveErrors = 0;
|
||||
const keys = Object.keys(output.hosts);
|
||||
|
||||
for (const [nodeKey, payload] of Object.entries(output.hosts)) {
|
||||
@@ -196,13 +248,19 @@ export async function createAgentManager(
|
||||
|
||||
log.info("agent", "Sync complete: %d nodes updated", keys.length);
|
||||
} catch (error) {
|
||||
consecutiveErrors++;
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
state.error = message;
|
||||
log.error("agent", "Sync failed: %s", message);
|
||||
log.error("agent", "Sync failed (%d/5): %s", consecutiveErrors, message);
|
||||
|
||||
if (consecutiveErrors >= 5) {
|
||||
log.warn("agent", "Too many consecutive failures, clearing state for next attempt");
|
||||
await rm(join(workDir, "tailscaled.state"), { force: true });
|
||||
}
|
||||
} finally {
|
||||
state.isSyncing = false;
|
||||
if (state.pendingResync) {
|
||||
state.pendingResync = false;
|
||||
isSyncing = false;
|
||||
if (pendingResync) {
|
||||
pendingResync = false;
|
||||
sync();
|
||||
}
|
||||
}
|
||||
@@ -299,7 +357,12 @@ export async function createAgentManager(
|
||||
},
|
||||
|
||||
dispose() {
|
||||
disposed = true;
|
||||
clearInterval(interval);
|
||||
if (proc) {
|
||||
proc.kill("SIGTERM");
|
||||
proc = null;
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
+36
-15
@@ -1,15 +1,27 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
|
||||
"github.com/tale/headplane/internal/config"
|
||||
"github.com/tale/headplane/internal/tsnet"
|
||||
"github.com/tale/headplane/internal/util"
|
||||
)
|
||||
|
||||
type output struct {
|
||||
Self string `json:"self"`
|
||||
Hosts map[string]json.RawMessage `json:"hosts"`
|
||||
}
|
||||
|
||||
type errorOutput struct {
|
||||
Error string `json:"error"`
|
||||
}
|
||||
|
||||
func main() {
|
||||
log := util.GetLogger()
|
||||
cfg, err := config.Load()
|
||||
@@ -23,21 +35,30 @@ func main() {
|
||||
|
||||
agent.Connect()
|
||||
|
||||
hosts, err := agent.FetchAllHostInfo(context.Background())
|
||||
if err != nil {
|
||||
log.Fatal("Failed to fetch host info: %s", err)
|
||||
}
|
||||
|
||||
output := struct {
|
||||
Self string `json:"self"`
|
||||
Hosts map[string]json.RawMessage `json:"hosts"`
|
||||
}{
|
||||
Self: agent.ID,
|
||||
Hosts: hosts,
|
||||
}
|
||||
|
||||
enc := json.NewEncoder(os.Stdout)
|
||||
if err := enc.Encode(output); err != nil {
|
||||
log.Fatal("Failed to encode result: %s", err)
|
||||
scanner := bufio.NewScanner(os.Stdin)
|
||||
|
||||
// Shut down cleanly on signal or stdin close
|
||||
sigCh := make(chan os.Signal, 1)
|
||||
signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT)
|
||||
|
||||
go func() {
|
||||
<-sigCh
|
||||
agent.Shutdown()
|
||||
os.Exit(0)
|
||||
}()
|
||||
|
||||
// Each line on stdin triggers a sync. The line content is ignored.
|
||||
for scanner.Scan() {
|
||||
hosts, err := agent.FetchAllHostInfo(context.Background())
|
||||
if err != nil {
|
||||
enc.Encode(errorOutput{Error: err.Error()})
|
||||
continue
|
||||
}
|
||||
|
||||
enc.Encode(output{
|
||||
Self: agent.ID,
|
||||
Hosts: hosts,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user