Something went wrong. Try again.
small gleam coding and (not yet) persistent agent daemon with a detachable cli
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510package daemon
import ( "bytes" "context" "crypto/rand" "encoding/hex" "encoding/json" "errors" "fmt" "io" "net/http" "os" "os/exec" "path/filepath" "regexp" "strconv" "strings" "time")
var sanitizerRegex = regexp.MustCompile(`[\p{Cc}\x{202a}-\x{202e}\x{2066}-\x{2069}]`)
// SessionText replaces control characters and direction overrides in session metadata.func SessionText(s string) string { return sanitizerRegex.ReplaceAllString(s, " ") }
const ( pollInterval = 100 * time.Millisecond pollAttempts = 300)
type Connection struct { Port int `json:"port"` Token string `json:"token"` Pid int `json:"pid"` Version int `json:"version"` Build string `json:"build,omitempty"`}
// Stale is a live daemon from a build other than the one this client bundles.type Stale struct { Running *Connection Bundled string}
// Replace decides whether a stale daemon stops so the bundled build can start.type Replace func(Stale) bool
type Session struct { ID string `json:"id"` Title string `json:"title,omitempty"` LastAssistantAt *int64 `json:"last_assistant_at,omitempty"` Workspace string `json:"workspace"` Model string `json:"model"` Protocol string `json:"protocol"` Provider string `json:"provider"`}
func AssistantAge(timestamp *int64, now time.Time) string { if timestamp == nil { return "time unknown" } sec := now.Unix() - *timestamp if sec < 0 { sec = 0 } if sec < 60 { return "just now" } if sec < 3600 { return fmt.Sprintf("%dm ago", sec/60) } if sec < 86400 { return fmt.Sprintf("%dh ago", sec/3600) } return fmt.Sprintf("%dd ago", sec/86400)}
func SessionListing(sessions []Session, now time.Time) string { if len(sessions) == 0 { return "no sessions" }
lines := make([]string, 0, len(sessions)) for _, s := range sessions { title := SessionText(s.Title) title = strings.TrimSpace(title) if title == "" { title = "session name unavailable" } shortID := s.ID if len(shortID) > 8 { shortID = shortID[:8] } lines = append(lines, fmt.Sprintf("%s [%s]\n last assistant: %s", title, shortID, AssistantAge(s.LastAssistantAt, now))) }
return strings.Join(lines, "\n\n")}
func Existing(homeDir string) (*Connection, error) { recordPath := filepath.Join(homeDir, "daemon.json") fi, err := os.Stat(recordPath) if err != nil || fi.Size() > 64*1024 { return nil, nil } data, err := os.ReadFile(recordPath) if err != nil { return nil, nil }
var conn Connection if err := json.Unmarshal(data, &conn); err != nil { return nil, nil }
if conn.Port < 1 || conn.Port > 65535 || conn.Token == "" || (conn.Version != 1 && conn.Version != 2) { return nil, nil }
ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond) defer cancel()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, fmt.Sprintf("http://127.0.0.1:%d/health", conn.Port), nil) if err != nil { return nil, nil } req.Header.Set("Authorization", "Bearer "+conn.Token)
client := &http.Client{ CheckRedirect: func(req *http.Request, via []*http.Request) error { return http.ErrUseLastResponse }, } res, err := client.Do(req) if err != nil { return nil, nil } defer res.Body.Close()
if res.StatusCode != http.StatusOK { return nil, nil }
body, err := readBounded(res.Body, 64*1024) if err != nil { return nil, nil }
var health struct { Version int `json:"version"` } if err := json.Unmarshal(body, &health); err != nil { return nil, nil }
if health.Version == conn.Version { return &conn, nil }
return nil, nil}
func checkCompatible(conn *Connection) (*Connection, error) { if conn.Version != 2 { return nil, errors.New("an older daemon is running; when its work is finished, run albedo daemon --stop, then start albedo again") } return conn, nil}
func resolveDaemonExecutable(path string) (string, error) { if !filepath.IsAbs(path) { return "", fmt.Errorf("ALBEDO_DAEMON must be an absolute executable path: %q", path) } resolved, err := exec.LookPath(path) if err != nil { return "", fmt.Errorf("invalid ALBEDO_DAEMON executable %q: %w", path, err) } fi, err := os.Stat(resolved) if err != nil { return "", fmt.Errorf("invalid ALBEDO_DAEMON executable %q: %w", path, err) } if fi.IsDir() { return "", fmt.Errorf("invalid ALBEDO_DAEMON executable %q: is a directory", path) } return resolved, nil}
// bundledBuild names a packaged daemon by its resolved executable, which a// content-addressed install changes on every build. A source checkout has none.func bundledBuild(daemonExe string) string { if daemonExe == "" { return "" } resolved, err := resolveDaemonExecutable(daemonExe) if err != nil { return "" } if real, err := filepath.EvalSymlinks(resolved); err == nil { return real } return resolved}
func buildDaemonEnv(homeDir, tokenHex, build string) []string { var env []string filteredVars := map[string]bool{ "ALBEDO_API_KEY": true, "ALBEDO_MODEL": true, "ALBEDO_BASE_URL": true, "ALBEDO_PROTOCOL": true, "ALBEDO_BUILD": true, } erlFlags := daemonErlFlags for _, kv := range os.Environ() { k, v, _ := strings.Cut(kv, "=") if k == "ERL_FLAGS" { // The operator's flags come last so they override the defaults. erlFlags += " " + v continue } if !filteredVars[k] { env = append(env, kv) } } env = append(env, "ERL_FLAGS="+erlFlags, "ALBEDO_HOME="+homeDir, "ALBEDO_TOKEN="+tokenHex) if build != "" { env = append(env, "ALBEDO_BUILD="+build) } return env}
// Sized for one local daemon, measured with ALBEDO_INSPECT (three large// sessions running turns at once peaked near 80 MB instead of 100).//// +P/+Q: the process and port tables are preallocated at their limits; the// defaults (1,048,576 processes, 65,536 ports) cost about 16 MB up front.//// +MB*/+MH* (binary and heap allocators): transcripts are loaded and dropped// per turn, so allocations come in bursts. Small carriers, address-order// best fit, a low single-block threshold (so large binaries get their own// mapping) and no carrier pooling let freed bursts go back to the OS instead// of staying resident as empty carrier space.const daemonErlFlags = "+P 65536 +Q 16384" + " +MBsbct 16 +MHsbct 32 +MBlmbcs 256 +MHlmbcs 256 +MBsmbcs 32 +MHsmbcs 32" + " +MBas aobf +MHas aobf +MBacul 0 +MHacul 0"
func daemonCommand(daemonExe, projectRoot string, env []string) (*exec.Cmd, error) { if daemonExe != "" { resolved, err := resolveDaemonExecutable(daemonExe) if err != nil { return nil, err } cmd := exec.Command(resolved) cmd.Dir = "" cmd.Env = env return cmd, nil }
cmd := exec.Command("gleam", "run") cmd.Dir = projectRoot cmd.Env = env return cmd, nil}
// stop shuts a daemon down and waits until its process has exited, so the// next daemon can take the home.func stop(homeDir string, conn *Connection) error { _, _ = Request[any](context.Background(), conn, "/shutdown", map[string]any{}) for attempt := 0; attempt < pollAttempts; attempt++ { running, err := Existing(homeDir) if err != nil { return err } if running == nil && !processAlive(conn.Pid) { return nil } time.Sleep(pollInterval) } return fmt.Errorf("daemon %d did not exit; inspect %s/daemon.log", conn.Pid, homeDir)}
// Ensure connects to the running daemon or starts one. When this client bundles// a daemon and the running one is another build, replace (if non-nil) decides// whether it is stopped first.func Ensure(homeDir, projectRoot string, replace Replace) (*Connection, error) { daemonExe := os.Getenv("ALBEDO_DAEMON") build := bundledBuild(daemonExe)
current, err := Existing(homeDir) if err != nil { return nil, err } if current != nil { if build == "" || current.Build == build || replace == nil || !replace(Stale{current, build}) { return checkCompatible(current) } if err := stop(homeDir, current); err != nil { return nil, err } }
if daemonExe != "" { if _, err := resolveDaemonExecutable(daemonExe); err != nil { return nil, err } }
if err := os.MkdirAll(homeDir, 0700); err != nil { return nil, err } _ = os.Chmod(homeDir, 0700)
lockPath := filepath.Join(homeDir, "starting.lock") lockFile, err := os.OpenFile(lockPath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0600) if err != nil && os.IsExist(err) && staleLock(lockPath) { _ = os.Remove(lockPath) lockFile, err = os.OpenFile(lockPath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0600) } if err != nil { if os.IsExist(err) { for attempt := 0; attempt < pollAttempts; attempt++ { running, err := Existing(homeDir) if err != nil { return nil, err } if running != nil { return checkCompatible(running) } time.Sleep(pollInterval) } return nil, fmt.Errorf("daemon startup timed out; inspect %s/daemon.log; remove %s if its starter is no longer running", homeDir, lockPath) } return nil, err } _, _ = fmt.Fprintf(lockFile, "%d", os.Getpid()) defer func() { _ = lockFile.Close() _ = os.Remove(lockPath) }()
again, err := Existing(homeDir) if err != nil { return nil, err } if again != nil { return checkCompatible(again) }
logPath := filepath.Join(homeDir, "daemon.log") logFile, err := os.OpenFile(logPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0600) if err != nil { return nil, err }
randomToken := make([]byte, 32) _, _ = rand.Read(randomToken) tokenHex := hex.EncodeToString(randomToken)
env := buildDaemonEnv(homeDir, tokenHex, build) cmd, err := daemonCommand(daemonExe, projectRoot, env) if err != nil { _ = logFile.Close() return nil, err } cmd.Stdout = logFile cmd.Stderr = logFile detach(cmd)
logStart, _ := logFile.Seek(0, io.SeekEnd) if err := cmd.Start(); err != nil { _ = logFile.Close() return nil, err } _ = logFile.Close()
exited := make(chan error, 1) go func() { exited <- cmd.Wait() }()
for attempt := 0; attempt < pollAttempts; attempt++ { running, err := Existing(homeDir) if err != nil { return nil, err } if running != nil { return checkCompatible(running) } select { case waitErr := <-exited: root := projectRoot if daemonExe != "" { root = "" } return nil, startupExitError(waitErr, logPath, logStart, root) case <-time.After(pollInterval): } }
return nil, fmt.Errorf("daemon startup timed out; inspect %s/daemon.log", homeDir)}
// staleLock reports whether the startup lock was left by a starter that is no// longer running (e.g. one interrupted with ctrl-c before it could clean up).func staleLock(lockPath string) bool { data, err := os.ReadFile(lockPath) if err != nil { return false } pid, err := strconv.Atoi(strings.TrimSpace(string(data))) if err != nil { // Locks from older clients carry no pid; treat them as stale once they // are older than the startup window. fi, statErr := os.Stat(lockPath) return statErr == nil && time.Since(fi.ModTime()) > pollInterval*pollAttempts } return !processAlive(pid)}
func startupExitError(waitErr error, logPath string, logStart int64, projectRoot string) error { msg := "daemon exited during startup" if waitErr != nil { msg = fmt.Sprintf("%s (%v)", msg, waitErr) } if f, err := os.Open(logPath); err == nil { defer f.Close() if _, err := f.Seek(logStart, io.SeekStart); err == nil { out, _ := io.ReadAll(io.LimitReader(f, 4096)) if tail := strings.TrimSpace(string(out)); tail != "" { msg += ":\n" + tail } } } if projectRoot != "" { msg += fmt.Sprintf("\n(project root: %s; set ALBEDO_ROOT or install with -ldflags \"-X main.buildRoot=...\")", projectRoot) } return errors.New(msg)}
func Request[T any](ctx context.Context, conn *Connection, path string, body any) (T, error) { method := http.MethodGet if body != nil { method = http.MethodPost } return RequestMethod[T](ctx, conn, method, path, body)}
func RequestMethod[T any](ctx context.Context, conn *Connection, method, path string, body any) (T, error) { var zero T if ctx == nil { ctx = context.Background() }
reqCtx, cancel := context.WithTimeout(ctx, 20*time.Second) defer cancel()
var bodyReader io.Reader if body != nil { encoded, err := json.Marshal(body) if err != nil { return zero, err } bodyReader = bytes.NewReader(encoded) }
reqURL := fmt.Sprintf("http://127.0.0.1:%d%s", conn.Port, path) req, err := http.NewRequestWithContext(reqCtx, method, reqURL, bodyReader) if err != nil { return zero, err }
req.Header.Set("Authorization", "Bearer "+conn.Token) req.Header.Set("Content-Type", "application/json")
client := &http.Client{ CheckRedirect: func(req *http.Request, via []*http.Request) error { return http.ErrUseLastResponse }, }
res, err := client.Do(req) if err != nil { return zero, err } defer res.Body.Close()
respData, err := readBounded(res.Body, 50*1024*1024) if err != nil { return zero, err }
if res.StatusCode < 200 || res.StatusCode >= 300 { var errResp struct { Error string `json:"error"` } _ = json.Unmarshal(respData, &errResp) if errResp.Error != "" { return zero, errors.New(errResp.Error) } return zero, fmt.Errorf("HTTP %d", res.StatusCode) }
var result T if len(respData) > 0 { if err := json.Unmarshal(respData, &result); err != nil { return zero, err } } return result, nil}