Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 4 additions & 17 deletions commands/agents.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,19 +66,6 @@ var (
colMuted = charm.Colors.Muted
)

// Stream-state transport frames (SSE kind "stream.state"). Defined locally so
// doctl builds against godo pins that have not yet exported HostedAgentEventKindStreamState
// / HostedAgentStreamState*. Wire values match the published godo API.
const (
hostedAgentEventKindStreamState godo.HostedAgentEventKind = "stream.state"
hostedAgentStreamStateSuperseded = "superseded"
)

type hostedAgentStreamState struct {
State string `json:"state"`
Cursor string `json:"cursor,omitempty"`
}

// detectStyling reports whether ANSI styling should be emitted for the current
// process: stdout is a terminal and NO_COLOR is unset.
func detectStyling() bool {
Expand Down Expand Up @@ -1148,7 +1135,7 @@ func RunAgentsLogs(c *CmdConfig) error {
for stream.Next() {
ev := stream.Current()
// Connection health, not session activity — never part of the history.
if ev.Kind == hostedAgentEventKindStreamState {
if ev.Kind == godo.HostedAgentEventKindStreamState {
continue
}
if ev.Kind == godo.HostedAgentEventKindTokenChunk {
Expand Down Expand Up @@ -1544,9 +1531,9 @@ func drainStream(stream *godo.HostedAgentSessionStream, out io.Writer, pending *

// stream.state reports the health of the connection, not session
// activity, so it never renders and never moves the cursor.
if ev.Kind == hostedAgentEventKindStreamState {
var st hostedAgentStreamState
if err := json.Unmarshal(ev.Payload, &st); err == nil && st.State == hostedAgentStreamStateSuperseded {
if ev.Kind == godo.HostedAgentEventKindStreamState {
var st godo.HostedAgentStreamState
if err := json.Unmarshal(ev.Payload, &st); err == nil && st.State == godo.HostedAgentStreamStateSuperseded {
thinking.stop()
acc.flush(out)
flushAwaitingApproval(out, &awaiting)
Expand Down
11 changes: 6 additions & 5 deletions commands/agents_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2182,7 +2182,7 @@ func TestStreamWithReconnect_supersededStopsWithoutReconnect(t *testing.T) {
stubReconnectSleep(t)

body := sseFrame("evt-1", string(godo.HostedAgentEventKindSessionUpdated), `{}`) +
sseFrame("", string(hostedAgentEventKindStreamState), `{"state":"superseded","cursor":""}`)
sseFrame("", string(godo.HostedAgentEventKindStreamState), `{"state":"superseded","cursor":""}`)
srv := httptest.NewServer(hostedAgentSSEHandler(body, nil))
t.Cleanup(srv.Close)

Expand Down Expand Up @@ -2241,9 +2241,9 @@ func TestDrainStream_HITLReattachShowsCommand(t *testing.T) {
// frame is transport bookkeeping: it renders nothing and must not become the
// reconnect cursor, or a reconnect would resume from a position no event holds.
func TestDrainStream_skipsStreamStateControlFrames(t *testing.T) {
body := sseFrame("", string(hostedAgentEventKindStreamState), `{"state":"live","cursor":""}`) +
body := sseFrame("", string(godo.HostedAgentEventKindStreamState), `{"state":"live","cursor":""}`) +
sseFrame("evt-7", string(godo.HostedAgentEventKindSessionUpdated), `{}`) +
sseFrame("", string(hostedAgentEventKindStreamState), `{"state":"catching_up","cursor":""}`)
sseFrame("", string(godo.HostedAgentEventKindStreamState), `{"state":"catching_up","cursor":""}`)
srv := httptest.NewServer(hostedAgentSSEHandler(body, nil))
t.Cleanup(srv.Close)

Expand Down Expand Up @@ -2358,8 +2358,9 @@ func TestStreamWithReconnect_replayCursorAfterMidStreamDrop(t *testing.T) {
mu.Lock()
calls++
n := calls
// Resume cursor rides as replay_from on control-plane /stream.
replayFrom := r.URL.Query().Get("replay_from")
// The live stream carries the resume cursor in the standard SSE
// Last-Event-ID header, not a replay_from query parameter.
replayFrom := r.Header.Get("Last-Event-ID")
mu.Unlock()

w.Header().Set("Content-Type", "text/event-stream")
Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ go 1.25.0
require (
github.com/blang/semver v3.5.1+incompatible
github.com/creack/pty v1.1.21
github.com/digitalocean/godo v1.202.0-beta.1
github.com/digitalocean/godo v1.202.1-0.20260731143800-4a1a63adeadb
github.com/docker/cli v24.0.5+incompatible
github.com/docker/docker v25.0.6+incompatible
github.com/docker/docker-credential-helpers v0.7.0 // indirect
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -120,8 +120,8 @@ github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSs
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/digitalocean/godo v1.202.0-beta.1 h1:1+wJiSwcshFgLdDHCr9MnGlnFUihox7EpY8cFxf76VA=
github.com/digitalocean/godo v1.202.0-beta.1/go.mod h1:xQsWpVCCbkDrWisHA72hPzPlnC+4W5w/McZY5ij9uvU=
github.com/digitalocean/godo v1.202.1-0.20260731143800-4a1a63adeadb h1:rBMYNfU4K4N8Uz5hiLTG2DId6k/TYA/ETsFvW9USJ2k=
github.com/digitalocean/godo v1.202.1-0.20260731143800-4a1a63adeadb/go.mod h1:xQsWpVCCbkDrWisHA72hPzPlnC+4W5w/McZY5ij9uvU=
github.com/distribution/reference v0.6.0 h1:0IXCQ5g4/QMHHkarYzh5l+u8T3t73zM5QvfrDyIgxBk=
github.com/distribution/reference v0.6.0/go.mod h1:BbU0aIcezP1/5jX/8MP0YiH4SdvB5Y4f/wlDRiLyi3E=
github.com/dlclark/regexp2 v1.11.5 h1:Q/sSnsKerHeCkc/jSTNq1oCm7KiVgUMZRDUoRu0JQZQ=
Expand Down
25 changes: 15 additions & 10 deletions internal/agentproxy/agentproxytest/harness.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,12 @@ import (
"net/http/httptest"
"sync"
"testing"

"github.com/digitalocean/godo"
)

// Event is one canned SSE event the harness streams back from
// GET /v2/agents/sessions/{id}/stream, matching the event-specific part of
// GET /v2/agents/sessions/{id}/events, matching the event-specific part of
// godo.HostedAgentEvent's wire shape (see HostedAgentEventKind's doc comment
// for the canonical type strings).
type Event struct {
Expand Down Expand Up @@ -57,7 +59,7 @@ type Harness struct {
mu sync.Mutex
sessionID string
runID string // returned by the next POST .../input call
events []Event // streamed, in order, by the next GET .../stream call
events []Event // streamed, in order, by the next GET .../events call
}

// New starts the fake harness and registers its shutdown via t.Cleanup.
Expand All @@ -69,8 +71,11 @@ func New(t *testing.T, sessionID string) *Harness {

mux := http.NewServeMux()
mux.HandleFunc("GET /v2/agents/sessions/{id}", h.handleGetSession)
// Temporarily on control-plane .../stream until OHP /events is on stage2.
mux.HandleFunc("GET /v2/agents/sessions/{id}/stream", h.handleStream)
// Live streaming is served by the data plane at .../events. The control
// plane's .../stream is deliberately not registered: it serves only
// replay-only reads, which no agentproxy caller makes, so a request landing
// there is a bug worth failing on.
mux.HandleFunc("GET /v2/agents/sessions/{id}/events", h.handleStream)
mux.HandleFunc("POST /v2/agents/sessions/{id}/input", h.handleInput)
mux.HandleFunc("POST /v2/agents/sessions/{id}/hitl/{requestID}", h.handleHITL)

Expand All @@ -80,7 +85,7 @@ func New(t *testing.T, sessionID string) *Harness {
}

// QueueRun arranges for the next POST .../input call to return runID, and
// for GET .../stream to then emit events (in order, tagged with runID),
// for GET .../events to then emit events (in order, tagged with runID),
// flushing after each one so a concurrent reader observes them incrementally
// rather than all at once at EOF.
//
Expand Down Expand Up @@ -139,20 +144,20 @@ func (h *Harness) handleStream(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
flusher, canFlush := w.(http.Flusher)

// Optional stream.state control frame (same wire value as the data-plane
// transport). Consumers must skip it — emit it so tests exercise that.
const streamStateKind = "stream.state"
// The data plane opens every stream with a stream.state control frame. It
// belongs to no run, so consumers must skip it rather than mistake it for
// session activity — emit it here so tests exercise that.
streamState, err := json.Marshal(eventWire{
TenantID: "15726539",
SessionID: sessionID,
Timestamp: "2026-01-01T00:00:00Z",
Type: streamStateKind,
Type: string(godo.HostedAgentEventKindStreamState),
Data: json.RawMessage(`{"state":"live","cursor":""}`),
})
if err != nil {
panic(fmt.Sprintf("agentproxytest: stream.state does not marshal to JSON: %v", err))
}
fmt.Fprintf(w, "event: %s\ndata: %s\n\n", streamStateKind, streamState)
fmt.Fprintf(w, "event: %s\ndata: %s\n\n", godo.HostedAgentEventKindStreamState, streamState)
if canFlush {
flusher.Flush()
}
Expand Down
89 changes: 78 additions & 11 deletions vendor/github.com/digitalocean/godo/hosted_agents.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion vendor/modules.txt
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,7 @@ github.com/creack/pty
# github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc
## explicit
github.com/davecgh/go-spew/spew
# github.com/digitalocean/godo v1.202.0-beta.1
# github.com/digitalocean/godo v1.202.1-0.20260731143800-4a1a63adeadb
## explicit; go 1.23.0
github.com/digitalocean/godo
github.com/digitalocean/godo/metrics
Expand Down