Something went wrong. Try again.
small gleam coding and (not yet) persistent agent daemon with a detachable cli
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885package daemon
import ( "bufio" "bytes" "context" "crypto/rand" "encoding/hex" "encoding/json" "errors" "fmt" "io" "net/http" "net/url" "strconv" "strings" "sync" "time")
type AgentPhase string
const ( PhaseResting AgentPhase = "resting" PhasePreparing AgentPhase = "preparing" PhaseReasoning AgentPhase = "reasoning" PhaseTool AgentPhase = "tool" PhaseWaiting AgentPhase = "waiting" PhaseCompacting AgentPhase = "compacting")
var validAgentPhases = map[AgentPhase]bool{ PhaseResting: true, PhasePreparing: true, PhaseReasoning: true, PhaseTool: true, PhaseWaiting: true, PhaseCompacting: true,}
type AgentStatus struct { Running bool `json:"running"` Idle bool `json:"idle"` Phase *AgentPhase `json:"phase,omitempty"`}
type WorkspaceUpdate struct { Workspace string `json:"workspace"`}
type WorkspaceMissingError struct { Workspace string}
func (e *WorkspaceMissingError) Error() string { return "workspace not found: " + e.Workspace}
type SendResult struct { OK bool `json:"ok"` Queued bool `json:"queued,omitempty"`}
type ChatClientOptions struct { BaseURL string Token string AgentID string ClientID string HTTPClient *http.Client}
type ChatClient struct { baseURL string token string agentID string clientID string httpClient *http.Client mu sync.Mutex afterSeq int argumentsByCall map[string]string}
func readBounded(r io.Reader, limit int64) ([]byte, error) { data, err := io.ReadAll(io.LimitReader(r, limit+1)) if err != nil { return nil, err } if int64(len(data)) > limit { return nil, fmt.Errorf("response body exceeded limit of %d bytes", limit) } return data, nil}
func NewChatClient(opts ChatClientOptions) *ChatClient { baseURL := strings.TrimRight(opts.BaseURL, "/") clientID := opts.ClientID if clientID == "" { b := make([]byte, 16) _, _ = rand.Read(b) clientID = "cli-" + hex.EncodeToString(b) } httpClient := opts.HTTPClient if httpClient == nil { httpClient = &http.Client{ CheckRedirect: func(req *http.Request, via []*http.Request) error { return http.ErrUseLastResponse }, } } return &ChatClient{ baseURL: baseURL, token: opts.Token, agentID: strings.TrimSpace(opts.AgentID), clientID: clientID, httpClient: httpClient, afterSeq: -1, argumentsByCall: make(map[string]string), }}
func (c *ChatClient) ClientID() string { return c.clientID}
func (c *ChatClient) agentURL(path string) string { if c.agentID != "" { return fmt.Sprintf("%s/sessions/%s%s", c.baseURL, url.PathEscape(c.agentID), path) } return c.baseURL + path}
func (c *ChatClient) parseResponseError(res *http.Response) error { body, err := readBounded(res.Body, 64*1024) if err != nil { return fmt.Errorf("%d %s: %w", res.StatusCode, res.Status, err) }
var data struct { Code string `json:"code"` Error string `json:"error"` Workspace string `json:"workspace"` } if err := json.Unmarshal(body, &data); err == nil { if data.Code == "workspace_missing" && data.Workspace != "" { return &WorkspaceMissingError{Workspace: data.Workspace} } if data.Error != "" { return errors.New(data.Error) } }
fallback := strings.TrimSpace(fmt.Sprintf("%d %s", res.StatusCode, res.Status)) if fallback == "" { fallback = "request failed" } return errors.New(fallback)}
func (c *ChatClient) Send(ctx context.Context, content string, image *ImageAttachment) (*SendResult, error) { reqCtx, cancel := context.WithTimeout(ctx, 20*time.Second) defer cancel()
payload := map[string]any{ "content": content, } if image != nil { payload["image"] = image } if c.clientID != "" { payload["clientId"] = c.clientID }
bodyData, err := json.Marshal(payload) if err != nil { return nil, err }
route := "/trigger/chat" if c.agentID != "" { route = "/events" }
req, err := http.NewRequestWithContext(reqCtx, http.MethodPost, c.agentURL(route), bytes.NewReader(bodyData)) if err != nil { return nil, err } req.Header.Set("Content-Type", "application/json") if c.token != "" { req.Header.Set("Authorization", "Bearer "+c.token) }
res, err := c.httpClient.Do(req) if err != nil { return nil, err } defer res.Body.Close()
if res.StatusCode < 200 || res.StatusCode >= 300 { return nil, c.parseResponseError(res) }
body, err := readBounded(res.Body, 64*1024) if err != nil { return nil, err }
var data struct { OK bool `json:"ok"` Queued bool `json:"queued"` } _ = json.Unmarshal(body, &data)
return &SendResult{ OK: true, Queued: data.Queued, }, nil}
func (c *ChatClient) ReplaceWorkspace(ctx context.Context, workspace string) (*WorkspaceUpdate, error) { if c.agentID == "" { return nil, errors.New("cannot replace workspace without agent ID") }
reqCtx, cancel := context.WithTimeout(ctx, 20*time.Second) defer cancel()
healthReq, err := http.NewRequestWithContext(reqCtx, http.MethodGet, c.baseURL+"/health", nil) if err != nil { return nil, err } if c.token != "" { healthReq.Header.Set("Authorization", "Bearer "+c.token) }
healthRes, err := c.httpClient.Do(healthReq) if err != nil { return nil, err } defer healthRes.Body.Close()
if healthRes.StatusCode != http.StatusOK { return nil, c.parseResponseError(healthRes) }
healthBody, err := readBounded(healthRes.Body, 64*1024) if err != nil { return nil, err }
var health struct { Capabilities []string `json:"capabilities"` } if err := json.Unmarshal(healthBody, &health); err != nil { return nil, errors.New("daemon returned invalid session metadata") }
hasWorkspaceCap := false for _, cap := range health.Capabilities { if cap == "session_workspace" { hasWorkspaceCap = true break } } if !hasWorkspaceCap { return nil, errors.New("daemon upgrade needed to change this workspace; when ready, run albedo daemon --stop, then albedo (this clears python variables)") }
payload := map[string]string{ "workspace": workspace, } bodyData, err := json.Marshal(payload) if err != nil { return nil, err }
req, err := http.NewRequestWithContext(reqCtx, http.MethodPost, c.agentURL("/workspace"), bytes.NewReader(bodyData)) if err != nil { return nil, err } req.Header.Set("Content-Type", "application/json") if c.token != "" { req.Header.Set("Authorization", "Bearer "+c.token) }
res, err := c.httpClient.Do(req) if err != nil { return nil, err } defer res.Body.Close()
if res.StatusCode != http.StatusOK { return nil, c.parseResponseError(res) }
wsBody, err := readBounded(res.Body, 64*1024) if err != nil { return nil, err }
var data struct { Workspace string `json:"workspace"` } if err := json.Unmarshal(wsBody, &data); err != nil || data.Workspace == "" { return nil, errors.New("daemon returned invalid session metadata") }
return &WorkspaceUpdate{Workspace: data.Workspace}, nil}
func (c *ChatClient) Interrupt(ctx context.Context) (bool, error) { reqCtx, cancel := context.WithTimeout(ctx, 15*time.Second) defer cancel()
route := "/awp/interrupt" if c.agentID != "" { route = "/interrupt" }
req, err := http.NewRequestWithContext(reqCtx, http.MethodPost, c.agentURL(route), bytes.NewReader([]byte("{}"))) if err != nil { return false, err } req.Header.Set("Content-Type", "application/json") if c.token != "" { req.Header.Set("Authorization", "Bearer "+c.token) }
res, err := c.httpClient.Do(req) if err != nil { return false, err } defer res.Body.Close()
if res.StatusCode != http.StatusOK { return false, c.parseResponseError(res) }
body, err := readBounded(res.Body, 64*1024) if err != nil { return false, err }
var data struct { Interrupted bool `json:"interrupted"` } _ = json.Unmarshal(body, &data) return data.Interrupted, nil}
// ContextWindow reads the model's context window from the session's last// prepared request. It is nil when neither a catalog nor the configuration// knows it, or when no request has been prepared yet.func (c *ChatClient) ContextWindow(ctx context.Context) (*int, error) { reqCtx, cancel := context.WithTimeout(ctx, 20*time.Second) defer cancel() req, err := http.NewRequestWithContext(reqCtx, http.MethodGet, c.agentURL("/context"), nil) if err != nil { return nil, err } if c.token != "" { req.Header.Set("Authorization", "Bearer "+c.token) } res, err := c.httpClient.Do(req) if err != nil { return nil, err } defer res.Body.Close() if res.StatusCode != http.StatusOK { return nil, c.parseResponseError(res) } body, err := readBounded(res.Body, 1024*1024) if err != nil { return nil, err } var data struct { Window *int `json:"context_window_tokens"` } if err := json.Unmarshal(body, &data); err != nil { return nil, err } if data.Window != nil && *data.Window <= 0 { return nil, nil } return data.Window, nil}
func (c *ChatClient) GetStatus(ctx context.Context) (*AgentStatus, error) { reqCtx, cancel := context.WithTimeout(ctx, 20*time.Second) defer cancel()
req, err := http.NewRequestWithContext(reqCtx, http.MethodGet, c.agentURL("/status"), nil) if err != nil { return nil, err } if c.token != "" { req.Header.Set("Authorization", "Bearer "+c.token) }
res, err := c.httpClient.Do(req) if err != nil { return nil, err } defer res.Body.Close()
if res.StatusCode != http.StatusOK { return nil, c.parseResponseError(res) }
body, err := readBounded(res.Body, 64*1024) if err != nil { return nil, err }
var data struct { Running bool `json:"running"` Idle bool `json:"idle"` Phase *string `json:"phase"` } if err := json.Unmarshal(body, &data); err != nil { return nil, err }
var phase *AgentPhase if data.Phase != nil { p := AgentPhase(*data.Phase) if validAgentPhases[p] { phase = &p } }
return &AgentStatus{ Running: data.Running, Idle: data.Idle, Phase: phase, }, nil}
func parseStreamEvent(raw map[string]any) *StreamEvent { typ, _ := raw["type"].(string) if typ == "" { if txt, ok := raw["text"].(string); ok { return &StreamEvent{Type: EventText, Text: txt} } return nil }
var timestamp *int64 if ts, ok := raw["timestamp"].(float64); ok && ts >= 0 { t := int64(ts) timestamp = &t }
switch typ { case "tool_progress": progVal := raw["progress"] if progVal == nil { return &StreamEvent{Type: EventToolProgress, Progress: nil} } progMap, ok := progVal.(map[string]any) if !ok { return nil } callID, _ := progMap["callId"].(string) name, _ := progMap["name"].(string) phase, _ := progMap["phase"].(string) if len(callID) > 200 || len(name) > 100 || (phase != "generating" && phase != "running") { return nil }
var intent *ToolIntent if im, ok := progMap["intent"].(map[string]any); ok { k, _ := im["kind"].(string) t, _ := im["target"].(string) if (k == "write" || k == "edit" || k == "read" || k == "run") && len(t) <= 300 { intent = &ToolIntent{Kind: k, Target: t} } }
var code *ToolCodePreview if cm, ok := progMap["code"].(map[string]any); ok && phase == "generating" { off, hasOff := cm["offset"].(float64) txt, hasTxt := cm["text"].(string) if hasOff && hasTxt && off >= 0 && len(txt) <= 512 { code = &ToolCodePreview{Offset: int(off), Text: txt} } }
return &StreamEvent{ Type: EventToolProgress, Progress: &ToolProgress{ CallID: callID, Name: cleanLabel(name), Phase: phase, Intent: intent, Code: code, }, }
case "message": txt, _ := raw["text"].(string) return &StreamEvent{ Type: EventMessage, Role: "assistant", Text: txt, Timestamp: timestamp, }
case "interrupted": return &StreamEvent{Type: EventInterrupted}
case "retry": return &StreamEvent{Type: EventRetry}
case "reset": return &StreamEvent{Type: EventReset}
case "note": txt, _ := raw["text"].(string) return &StreamEvent{Type: EventNote, Text: txt}
case "user": txt, _ := raw["text"].(string) source, _ := raw["source"].(string) trig, _ := raw["triggeredAt"].(string) clientID, _ := raw["clientId"].(string) img := ParseImageMetadata(raw["image"]) return &StreamEvent{ Type: EventUser, Text: txt, Source: source, TriggeredAt: trig, ClientID: clientID, Timestamp: timestamp, Image: img, }
case "thinking": txt, _ := raw["text"].(string) return &StreamEvent{Type: EventThinking, Text: txt}
case "error": txt, _ := raw["text"].(string) return &StreamEvent{Type: EventError, Text: txt}
case "tool": name, _ := raw["name"].(string) result, _ := raw["result"].(string) args, _ := raw["args"].(map[string]any) trace := ParseToolTrace(raw["trace"]) return &StreamEvent{ Type: EventTool, ToolName: name, ToolArgs: args, ToolResult: result, ToolTrace: trace, }
case "usage": model, _ := raw["model"].(string) var recordedAt *int64 if ra, ok := raw["recordedAt"].(float64); ok { r := int64(ra) recordedAt = &r } getInt := func(k string) *int { if v, ok := raw[k].(float64); ok { i := int(v) return &i } return nil } getFloat := func(k string) *float64 { if v, ok := raw[k].(float64); ok { return &v } return nil } return &StreamEvent{ Type: EventUsage, Usage: &Usage{ Model: model, RecordedAt: recordedAt, PromptTokens: getInt("promptTokens"), CachedPromptTokens: getInt("cachedPromptTokens"), CacheWriteTokens: getInt("cacheWriteTokens"), CompletionTokens: getInt("completionTokens"), TotalTokens: getInt("totalTokens"), ElapsedMs: getFloat("elapsedMs"), TokensPerSecond: getFloat("tokensPerSecond"), }, }
case "compacted": evicted, hasEv := raw["evicted"].(float64) summary, hasSum := raw["summary"].(string) if hasEv && hasSum && evicted >= 0 && evicted <= 1000000 && len(summary) <= 60000 { return &StreamEvent{ Type: EventCompacted, Evicted: int(evicted), Summary: summary, } } return nil
default: if txt, ok := raw["text"].(string); ok { return &StreamEvent{Type: EventText, Text: txt} } return nil }}
func (c *ChatClient) Stream(ctx context.Context, onOpen func(), onEvent func(StreamEvent) error) error { c.mu.Lock() afterSeq := c.afterSeq c.mu.Unlock()
var route string if c.agentID != "" { route = fmt.Sprintf("/stream?after_seq=%d", afterSeq) } else { if afterSeq > 0 { route = fmt.Sprintf("/awp/stream?after_seq=%d", afterSeq) } else { route = "/awp/stream?tail=true" } }
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.agentURL(route), nil) if err != nil { return err } req.Header.Set("Accept", "text/event-stream") if c.token != "" { req.Header.Set("Authorization", "Bearer "+c.token) }
res, err := c.httpClient.Do(req) if err != nil { return err } defer res.Body.Close()
if res.StatusCode != http.StatusOK { return c.parseResponseError(res) }
if onOpen != nil { onOpen() }
reporter := NewToolProgressReporter(func(prog *ToolProgress) error { return onEvent(StreamEvent{Type: EventToolProgress, Progress: prog}) })
defer func() { _ = reporter.Report(nil, "running") if ctx.Err() != nil { c.mu.Lock() c.afterSeq = -1 c.argumentsByCall = make(map[string]string) c.mu.Unlock() } }()
scanner := bufio.NewScanner(res.Body) buf := make([]byte, 64*1024) scanner.Buffer(buf, 10*1024*1024)
var eventType string eventSeq := afterSeq
for scanner.Scan() { select { case <-ctx.Done(): return ctx.Err() default: }
line := scanner.Text() if line == "" { eventType = "" continue }
if strings.HasPrefix(line, "event:") { eventType = strings.TrimSpace(line[6:]) continue } if strings.HasPrefix(line, "id:") { if seq, err := strconv.Atoi(strings.TrimSpace(line[3:])); err == nil && seq > 0 { eventSeq = seq } continue } if !strings.HasPrefix(line, "data:") { continue }
payload := strings.TrimSpace(line[5:]) if payload == "" { continue }
var raw any if err := json.Unmarshal([]byte(payload), &raw); err != nil { continue }
if eventType == "error" { errStr := payload if m, ok := raw.(map[string]any); ok { if e, ok := m["error"].(string); ok { errStr = e } else if t, ok := m["text"].(string); ok { errStr = t } } return errors.New(errStr) }
rawMap, isMap := raw.(map[string]any) if isMap { if cursorVal, hasCursor := rawMap["cursor"]; hasCursor { if cursorNum, ok := cursorVal.(float64); ok { if eventsArr, ok := rawMap["events"].([]any); ok { // a reset opens the snapshot, and the snapshot is the rest of its page snapshot := false for _, evRaw := range eventsArr { evMap, ok := evRaw.(map[string]any) if !ok { continue } itemType, _ := evMap["type"].(string) switch itemType { case "reset", "retry", "turn_started", "message", "interrupted", "error": c.mu.Lock() c.argumentsByCall = make(map[string]string) c.mu.Unlock() if err := reporter.Report(nil, "running"); err != nil { return err } }
if itemType == "reset" { snapshot = true if err := onEvent(StreamEvent{Type: EventReset}); err != nil { return err } continue } if itemType == "turn_started" { continue } if itemType == "tool" { if callID, ok := evMap["callId"].(string); ok { c.mu.Lock() delete(c.argumentsByCall, callID) c.mu.Unlock() } if err := reporter.Report(nil, "running"); err != nil { return err } } if itemType == "tool_progress" { if prog, ok := evMap["progress"].(map[string]any); ok { ph, _ := prog["phase"].(string) cid, _ := prog["callId"].(string) if ph == "running" && cid != "" { c.mu.Lock() delete(c.argumentsByCall, cid) c.mu.Unlock() if err := reporter.Report(nil, "running"); err != nil { return err } } } } if itemType == "arguments_delta" { callID, _ := evMap["callId"].(string) text, _ := evMap["text"].(string) if callID != "" && text != "" { c.mu.Lock() fresh := false if _, exists := c.argumentsByCall[callID]; !exists { fresh = true }
retained := 0 for k, v := range c.argumentsByCall { retained += len(k) + len(v) }
extraCallID := 0 if fresh { extraCallID = len(callID) }
if (fresh && len(c.argumentsByCall) >= 32) || retained+len(text)+extraCallID > 2000000 { c.argumentsByCall = make(map[string]string) c.afterSeq = -1 c.mu.Unlock() return errors.New("tool argument previews exceed client limit") }
args := c.argumentsByCall[callID] + text c.argumentsByCall[callID] = args c.mu.Unlock()
call := &ToolCallAssembly{ ID: callID, } call.Function.Name = "python" call.Function.Arguments = args if err := reporter.Report(call, "generating"); err != nil { return err } continue } }
if itemType == "tool" { if strArgs, ok := evMap["args"].(string); ok { var parsedArgs map[string]any if err := json.Unmarshal([]byte(strArgs), &parsedArgs); err == nil { evMap["args"] = parsedArgs } else { evMap["args"] = make(map[string]any) } } }
if ev := parseStreamEvent(evMap); ev != nil { ev.Replayed = snapshot if err := onEvent(*ev); err != nil { return err } } }
c.mu.Lock() c.afterSeq = int(cursorNum) c.mu.Unlock() continue } } }
var targetMap map[string]any = rawMap if envType, ok := rawMap["type"].(string); ok { if envType == "stream.event" { if pMap, ok := rawMap["payload"].(map[string]any); ok { targetMap = pMap } } else if envType == "conversation.message" { if pMap, ok := rawMap["payload"].(map[string]any); ok { role, _ := pMap["role"].(string) content, _ := pMap["content"].(string) if role == "assistant" && strings.TrimSpace(content) != "" { targetMap = map[string]any{ "type": "message", "role": "assistant", "text": content, } } } } }
if ev := parseStreamEvent(targetMap); ev != nil { if err := onEvent(*ev); err != nil { return err } } c.mu.Lock() if eventSeq > c.afterSeq { c.afterSeq = eventSeq } c.mu.Unlock() } }
return scanner.Err()}