633 lines
23 KiB
Go
633 lines
23 KiB
Go
package httpserver_test
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"strings"
|
|
"sync/atomic"
|
|
"testing"
|
|
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
|
|
"github.com/krow/krow-backend/go-api/internal/gateway"
|
|
"github.com/krow/krow-backend/go-api/internal/httpserver"
|
|
"github.com/krow/krow-backend/go-api/internal/runtime"
|
|
"github.com/krow/krow-backend/go-api/internal/testutil"
|
|
"github.com/krow/krow-backend/go-api/internal/tools"
|
|
)
|
|
|
|
// The run endpoint's tests.
|
|
//
|
|
// Everything here is about the SEAM rather than the runtime — the runtime has
|
|
// its own tests and they are thorough. What this file asks is the set of
|
|
// questions only the HTTP layer can answer:
|
|
//
|
|
// - Does an unauthenticated caller get in?
|
|
// - Does another tenant's agent look absent or forbidden? (It must look
|
|
// absent — a 403 is a confirmation that the agent exists.)
|
|
// - Does a bounded run answer like a failure or like a run?
|
|
// - Does a pending confirmation reach the client in a shape it can act on?
|
|
// - Can one worker read another's trajectory?
|
|
//
|
|
// The model is scripted throughout. That is not a compromise: this file is
|
|
// about status codes and response shapes, and a live model would make it slow,
|
|
// non-deterministic and impossible to run without a credential.
|
|
|
|
/* ── Fixtures ───────────────────────────────────────────────────────────── */
|
|
|
|
// stubGateway answers with whatever it was given.
|
|
type stubGateway struct {
|
|
text string
|
|
calls []gateway.ToolCall
|
|
err error
|
|
sent int
|
|
}
|
|
|
|
func (s *stubGateway) Complete(_ context.Context, _ gateway.Request) (*gateway.Response, error) {
|
|
s.sent++
|
|
if s.err != nil {
|
|
return &gateway.Response{Model: "stub"}, s.err
|
|
}
|
|
if len(s.calls) > 0 && s.sent == 1 {
|
|
return &gateway.Response{
|
|
ToolCalls: s.calls, StopReason: "tool_use", Model: "stub",
|
|
Usage: gateway.Usage{InputTokens: 400, OutputTokens: 30},
|
|
}, nil
|
|
}
|
|
return &gateway.Response{
|
|
Text: s.text, StopReason: "end_turn", Model: "stub",
|
|
Usage: gateway.Usage{InputTokens: 500, OutputTokens: 40},
|
|
}, nil
|
|
}
|
|
|
|
// publishAgent writes a runnable agent definition.
|
|
func publishAgent(t *testing.T, pool *pgxpool.Pool, orgID, userID, id string, toolNames ...string) {
|
|
t.Helper()
|
|
var toolBlock string
|
|
if len(toolNames) > 0 {
|
|
toolBlock = "tools:\n"
|
|
for _, n := range toolNames {
|
|
toolBlock += " - " + n + "\n"
|
|
}
|
|
}
|
|
md := fmt.Sprintf(`---
|
|
id: %s
|
|
name: Test Agent
|
|
description: An agent for the run endpoint's tests
|
|
status: published
|
|
version: 1
|
|
pages:
|
|
- control-center
|
|
reasoning: balanced
|
|
%s---
|
|
|
|
## Instructions
|
|
Answer the question.
|
|
`, id, toolBlock)
|
|
|
|
if _, err := pool.Exec(context.Background(), `
|
|
INSERT INTO agent_definitions
|
|
(definition_id, org_id, visibility, created_by, markdown, status, version, name, description, pages)
|
|
VALUES ($1::text, $2::uuid, 'organization', $3::uuid, $4::text, 'published', 1,
|
|
'Test Agent', 'An agent for tests', ARRAY['control-center'])`,
|
|
id, orgID, userID, md); err != nil {
|
|
t.Fatalf("publish agent %s: %v", id, err)
|
|
}
|
|
}
|
|
|
|
// seedOrgAdmin creates a fresh tenant and an admin user inside it.
|
|
//
|
|
// A tenant per test, not the seeded one. The cross-tenant assertions below need
|
|
// two organizations that genuinely do not know about each other, and reusing
|
|
// the fixture's org for one of them would make "another tenant" mean "the same
|
|
// tenant with a different user".
|
|
func seedOrgAdmin(t *testing.T, h *testutil.Harness) (orgID, userID string) {
|
|
t.Helper()
|
|
slug := fmt.Sprintf("runs-%d-%s", orgCounter.Add(1), t.Name())
|
|
slug = strings.ToLower(strings.NewReplacer("/", "-", "_", "-", " ", "-").Replace(slug))
|
|
if len(slug) > 60 {
|
|
slug = slug[:60]
|
|
}
|
|
if err := h.Pool.QueryRow(context.Background(),
|
|
`INSERT INTO organizations (name, slug) VALUES ($1, $2) RETURNING id::text`,
|
|
slug, slug).Scan(&orgID); err != nil {
|
|
t.Fatalf("create org: %v", err)
|
|
}
|
|
userID = newUserWithRole(t, h.Pool, orgID,
|
|
fmt.Sprintf("owner-%s@runs.test", slug), "admin")
|
|
return orgID, userID
|
|
}
|
|
|
|
// orgCounter keeps fixture slugs unique. Emails and slugs are globally unique,
|
|
// so two tenants in one test collide without it.
|
|
var orgCounter atomic.Int64
|
|
|
|
// runServer builds a server whose runtime is driven by a scripted gateway.
|
|
func runServer(t *testing.T, h *testutil.Harness, gw gateway.Gateway, reg *tools.Registry) *httpserver.Server {
|
|
t.Helper()
|
|
engine := runtime.NewEngine(h.Pool, runtime.WithAgentExecutor(
|
|
runtime.NewModelExecutor(gw, runtime.NewPostgresSink(h.Pool), reg),
|
|
))
|
|
return newServer(t, h, nil, httpserver.WithAgentEngine(engine))
|
|
}
|
|
|
|
// postRun calls the run endpoint as one actor.
|
|
func postRun(t *testing.T, handler http.Handler, a actor, agentID, body string) (int, map[string]any) {
|
|
t.Helper()
|
|
req, err := http.NewRequest("POST",
|
|
"/api/v1/agents/"+agentID+"/runs", strings.NewReader(body))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
req.Header.Set("Content-Type", "application/json")
|
|
if a.cookie != nil {
|
|
req.AddCookie(a.cookie)
|
|
}
|
|
return doJSON(t, handler, req)
|
|
}
|
|
|
|
func doJSON(t *testing.T, handler http.Handler, req *http.Request) (int, map[string]any) {
|
|
t.Helper()
|
|
rec := httptest.NewRecorder()
|
|
handler.ServeHTTP(rec, req)
|
|
var body map[string]any
|
|
if rec.Body.Len() > 0 {
|
|
if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil {
|
|
t.Fatalf("response was not JSON: %s", rec.Body.String())
|
|
}
|
|
}
|
|
return rec.Code, body
|
|
}
|
|
|
|
/* ── The endpoint exists at all ─────────────────────────────────────────── */
|
|
|
|
func TestTheRunRoutesAreAbsentWithoutARuntime(t *testing.T) {
|
|
// A deployment with no model credential does not serve agents. Registering
|
|
// the routes anyway would accept runs and fail every one at the gateway —
|
|
// an outage shaped like a feature. 404 says "this deployment does not do
|
|
// that", which is true; 500 would say "this deployment is broken", which is
|
|
// not.
|
|
h := testutil.New(t)
|
|
srv := newServer(t, h, nil) // no WithAgentEngine, no API key
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "test-agent")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@runs.test", "admin")
|
|
|
|
code, _ := postRun(t, handler, admin, "test-agent", `{"input":"hello"}`)
|
|
if code != http.StatusNotFound {
|
|
t.Errorf("status %d without a runtime, want 404", code)
|
|
}
|
|
}
|
|
|
|
func TestAnUnauthenticatedRunIsRefused(t *testing.T) {
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &stubGateway{text: "hello"}, nil)
|
|
handler := srv.Handler()
|
|
|
|
code, _ := postRun(t, handler, actor{}, "test-agent", `{"input":"hello"}`)
|
|
if code != http.StatusUnauthorized {
|
|
t.Errorf("status %d for an unauthenticated run, want 401", code)
|
|
}
|
|
}
|
|
|
|
/* ── A completed run ────────────────────────────────────────────────────── */
|
|
|
|
func TestACompletedRunAnswersWithItsOutputAndCost(t *testing.T) {
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &stubGateway{text: "Twelve events, mostly logins."}, nil)
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "test-agent")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@runs.test", "admin")
|
|
|
|
code, body := postRun(t, handler, admin, "test-agent", `{"input":"what happened?"}`)
|
|
if code != http.StatusOK {
|
|
t.Fatalf("status %d: %v", code, body)
|
|
}
|
|
if body["termination"] != "Completed" {
|
|
t.Errorf("termination = %v, want Completed", body["termination"])
|
|
}
|
|
if body["output"] != "Twelve events, mostly logins." {
|
|
t.Errorf("output = %v", body["output"])
|
|
}
|
|
if body["runId"] == nil || body["runId"] == "" {
|
|
t.Error("a run came back with no id; nothing can point at its trajectory")
|
|
}
|
|
// Token accounting reaches the client. A caller paying for runs should be
|
|
// able to see what one cost without reading a log.
|
|
usage, _ := body["usage"].(map[string]any)
|
|
if usage == nil || usage["totalTokens"] == nil {
|
|
t.Errorf("no usage in the response: %v", body)
|
|
}
|
|
// A completed run carries no user-facing message: the output IS the answer.
|
|
if msg, ok := body["message"].(string); ok && msg != "" {
|
|
t.Errorf("a completed run carried a message: %q", msg)
|
|
}
|
|
}
|
|
|
|
/* ── The denial rules ───────────────────────────────────────────────────── */
|
|
|
|
func TestAnotherTenantsAgentIsAbsentRatherThanForbidden(t *testing.T) {
|
|
// §8's rule about denials applies to agents as much as to rows. If "exists
|
|
// but not yours" answered 403 and "no such agent" answered 404, the
|
|
// endpoint would be a way to enumerate other tenants' agents one id at a
|
|
// time — and the agent would happily run that enumeration.
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &stubGateway{text: "hello"}, nil)
|
|
handler := srv.Handler()
|
|
|
|
mine, mineAdmin := seedOrgAdmin(t, h)
|
|
theirs, theirsAdmin := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, theirs, theirsAdmin, "their-agent")
|
|
_ = mineAdmin
|
|
|
|
admin := signInAs(t, handler, h.Pool, mine, "admin", "admin@mine.test", "admin")
|
|
|
|
real, realBody := postRun(t, handler, admin, "their-agent", `{"input":"hi"}`)
|
|
fake, fakeBody := postRun(t, handler, admin, "no-such-agent-at-all", `{"input":"hi"}`)
|
|
|
|
if real != http.StatusNotFound {
|
|
t.Errorf("another tenant's agent answered %d, want 404", real)
|
|
}
|
|
if fake != http.StatusNotFound {
|
|
t.Errorf("an imaginary agent answered %d, want 404", fake)
|
|
}
|
|
if fmt.Sprint(realBody) != fmt.Sprint(fakeBody) {
|
|
t.Errorf("a real-but-forbidden agent is distinguishable from an imaginary one:\n"+
|
|
" theirs: %v\n invented: %v", realBody, fakeBody)
|
|
}
|
|
}
|
|
|
|
func TestARunWithNoInputIsRefused(t *testing.T) {
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &stubGateway{text: "hello"}, nil)
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "test-agent")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@runs.test", "admin")
|
|
|
|
code, _ := postRun(t, handler, admin, "test-agent", `{}`)
|
|
if code != http.StatusUnprocessableEntity {
|
|
t.Errorf("status %d for an empty input, want 422", code)
|
|
}
|
|
}
|
|
|
|
/* ── A pending confirmation ─────────────────────────────────────────────── */
|
|
|
|
func TestAPendingConfirmationReachesTheClientAsAQuestionNotAnError(t *testing.T) {
|
|
// I4 arriving at the surface. A run waiting on a person is not a failure:
|
|
// it has an id, a cost, a trajectory and a payload somebody has to read.
|
|
// Answering it 500 would make the whole write path look broken, and the
|
|
// client would have no token to call back with.
|
|
h := testutil.New(t)
|
|
|
|
var wrote int
|
|
reg := tools.NewRegistry()
|
|
reg.MustRegister(tools.Tool{
|
|
Name: "assign_worker", Description: "Assigns somebody to something, for this test.",
|
|
InputSchema: map[string]any{"type": "object"}, Effect: tools.EffectWrite,
|
|
Confirm: func(context.Context, tools.Context, json.RawMessage) (*tools.Confirmation, *tools.Result) {
|
|
return &tools.Confirmation{
|
|
Title: "Assign Maya Chen to Bar Supervisor",
|
|
Summary: "Maya Chen will be scheduled to work Friday evening.",
|
|
}, nil
|
|
},
|
|
Handler: func(context.Context, tools.Context, json.RawMessage) tools.Result {
|
|
wrote++
|
|
return tools.OK(map[string]any{"ok": true})
|
|
},
|
|
})
|
|
|
|
gw := &stubGateway{
|
|
text: "done",
|
|
calls: []gateway.ToolCall{{ID: "c1", Name: "assign_worker", Input: json.RawMessage(`{}`)}},
|
|
}
|
|
srv := runServer(t, h, gw, reg)
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "cover-agent", "assign_worker")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@runs.test", "admin")
|
|
|
|
code, body := postRun(t, handler, admin, "cover-agent", `{"input":"cover Friday"}`)
|
|
|
|
if code != http.StatusOK {
|
|
t.Fatalf("status %d for a pending confirmation, want 200: %v", code, body)
|
|
}
|
|
if body["termination"] != "ConfirmationPending" {
|
|
t.Fatalf("termination = %v, want ConfirmationPending", body["termination"])
|
|
}
|
|
if wrote != 0 {
|
|
t.Fatalf("the write ran %d times without an approval", wrote)
|
|
}
|
|
|
|
confirmations, _ := body["confirmations"].([]any)
|
|
if len(confirmations) != 1 {
|
|
t.Fatalf("%d confirmations in the response, want 1: %v", len(confirmations), body)
|
|
}
|
|
c, _ := confirmations[0].(map[string]any)
|
|
if c["token"] == nil || c["token"] == "" {
|
|
t.Error("the confirmation has no token; the client can never answer it")
|
|
}
|
|
if c["title"] == nil || c["title"] == "" {
|
|
t.Error("the confirmation has nothing written on it for a person to read")
|
|
}
|
|
// And the client is told what to say to the user, derived here rather than
|
|
// raised from the core.
|
|
if msg, _ := body["message"].(string); !strings.Contains(strings.ToLower(msg), "approve") {
|
|
t.Errorf("message = %q; it should tell the user an approval is needed", msg)
|
|
}
|
|
}
|
|
|
|
/* ── Reading a trajectory ───────────────────────────────────────────────── */
|
|
|
|
func TestATrajectoryIsReadableByItsOwnerAndNobodyElse(t *testing.T) {
|
|
// A trajectory holds the question that was asked and the records retrieved
|
|
// to answer it. "Anyone in the tenant may read any run" would let every
|
|
// worker read every colleague's conversation with an agent — including the
|
|
// ones about them.
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &stubGateway{text: "an answer"}, nil)
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "test-agent")
|
|
|
|
maya := signInAs(t, handler, h.Pool, orgID, "maya", "maya@runs.test", "talent")
|
|
dan := signInAs(t, handler, h.Pool, orgID, "dan", "dan@runs.test", "talent")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@runs.test", "admin")
|
|
|
|
code, body := postRun(t, handler, maya, "test-agent", `{"input":"my private question"}`)
|
|
if code != http.StatusOK {
|
|
t.Fatalf("status %d: %v", code, body)
|
|
}
|
|
runID, _ := body["runId"].(string)
|
|
if runID == "" {
|
|
t.Fatal("no run id came back")
|
|
}
|
|
|
|
get := func(a actor) (int, map[string]any) {
|
|
req, _ := http.NewRequest("GET", "/api/v1/runs/"+runID, nil)
|
|
if a.cookie != nil {
|
|
req.AddCookie(a.cookie)
|
|
}
|
|
return doJSON(t, handler, req)
|
|
}
|
|
|
|
if code, _ := get(maya); code != http.StatusOK {
|
|
t.Errorf("the owner could not read their own run: %d", code)
|
|
}
|
|
if code, _ := get(dan); code != http.StatusNotFound {
|
|
t.Errorf("another worker read a colleague's run: %d, want 404", code)
|
|
}
|
|
// An operator sees the organization's runs. That is what an operator
|
|
// console is, and it is the same reach the policy table already gives them
|
|
// over every other resource.
|
|
if code, _ := get(admin); code != http.StatusOK {
|
|
t.Errorf("an operator could not read their organization's run: %d", code)
|
|
}
|
|
}
|
|
|
|
func TestATrajectoryFromAnotherTenantIsAbsent(t *testing.T) {
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &stubGateway{text: "an answer"}, nil)
|
|
handler := srv.Handler()
|
|
|
|
mine, mineAdmin := seedOrgAdmin(t, h)
|
|
theirs, _ := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, mine, mineAdmin, "test-agent")
|
|
|
|
owner := signInAs(t, handler, h.Pool, mine, "owner", "owner@mine.test", "admin")
|
|
outsider := signInAs(t, handler, h.Pool, theirs, "outsider", "outsider@theirs.test", "admin")
|
|
|
|
_, body := postRun(t, handler, owner, "test-agent", `{"input":"a question"}`)
|
|
runID, _ := body["runId"].(string)
|
|
|
|
req, _ := http.NewRequest("GET", "/api/v1/runs/"+runID, nil)
|
|
req.AddCookie(outsider.cookie)
|
|
code, _ := doJSON(t, handler, req)
|
|
|
|
if code != http.StatusNotFound {
|
|
t.Errorf("another tenant read a run: %d, want 404", code)
|
|
}
|
|
}
|
|
|
|
/* ── Streaming ──────────────────────────────────────────────────────────── */
|
|
|
|
// streamingStub is a gateway that emits text in pieces.
|
|
type streamingStub struct {
|
|
pieces []string
|
|
deltas int
|
|
}
|
|
|
|
func (s *streamingStub) Complete(context.Context, gateway.Request) (*gateway.Response, error) {
|
|
return &gateway.Response{
|
|
Text: strings.Join(s.pieces, ""), StopReason: "end_turn", Model: "stub",
|
|
Usage: gateway.Usage{InputTokens: 100, OutputTokens: 20},
|
|
}, nil
|
|
}
|
|
|
|
func (s *streamingStub) Stream(_ context.Context, _ gateway.Request, onDelta func(string)) (*gateway.Response, error) {
|
|
for _, p := range s.pieces {
|
|
s.deltas++
|
|
onDelta(p)
|
|
}
|
|
return &gateway.Response{
|
|
Text: strings.Join(s.pieces, ""), StopReason: "end_turn", Model: "stub",
|
|
Usage: gateway.Usage{InputTokens: 100, OutputTokens: 20},
|
|
}, nil
|
|
}
|
|
|
|
// sseEvents pulls the JSON payloads out of an SSE body.
|
|
func sseEvents(t *testing.T, body string) []map[string]any {
|
|
t.Helper()
|
|
var out []map[string]any
|
|
for _, line := range strings.Split(body, "\n") {
|
|
line = strings.TrimSpace(line)
|
|
if !strings.HasPrefix(line, "data:") {
|
|
continue
|
|
}
|
|
payload := strings.TrimSpace(line[5:])
|
|
if payload == "" || payload == "[DONE]" {
|
|
continue
|
|
}
|
|
var e map[string]any
|
|
if err := json.Unmarshal([]byte(payload), &e); err != nil {
|
|
t.Fatalf("event was not JSON: %s", payload)
|
|
}
|
|
out = append(out, e)
|
|
}
|
|
return out
|
|
}
|
|
|
|
func TestAStreamedRunDeliversTextThenTheFinishedRun(t *testing.T) {
|
|
h := testutil.New(t)
|
|
gw := &streamingStub{pieces: []string{"Twelve ", "events, ", "mostly logins."}}
|
|
srv := runServer(t, h, gw, nil)
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "test-agent")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@stream.test", "admin")
|
|
|
|
req, _ := http.NewRequest("POST", "/api/v1/agents/test-agent/runs",
|
|
strings.NewReader(`{"input":"what happened?"}`))
|
|
req.Header.Set("Content-Type", "application/json")
|
|
req.Header.Set("Accept", "text/event-stream")
|
|
req.AddCookie(admin.cookie)
|
|
|
|
rec := httptest.NewRecorder()
|
|
handler.ServeHTTP(rec, req)
|
|
|
|
if rec.Code != http.StatusOK {
|
|
t.Fatalf("status %d: %s", rec.Code, rec.Body.String())
|
|
}
|
|
if ct := rec.Header().Get("Content-Type"); !strings.Contains(ct, "text/event-stream") {
|
|
t.Fatalf("Content-Type is %q, want an event stream — the response did not stream", ct)
|
|
}
|
|
|
|
events := sseEvents(t, rec.Body.String())
|
|
var deltas []string
|
|
var final map[string]any
|
|
for _, e := range events {
|
|
if d, ok := e["delta"].(string); ok {
|
|
deltas = append(deltas, d)
|
|
}
|
|
if r, ok := e["run"].(map[string]any); ok {
|
|
final = r
|
|
}
|
|
}
|
|
|
|
if len(deltas) != 3 {
|
|
t.Errorf("%d text deltas, want 3 — the text arrived in one piece", len(deltas))
|
|
}
|
|
if strings.Join(deltas, "") != "Twelve events, mostly logins." {
|
|
t.Errorf("the deltas do not reassemble into the answer: %q", strings.Join(deltas, ""))
|
|
}
|
|
|
|
// The property that keeps the two paths honest: a client that ignored every
|
|
// delta and read only the last event is where it would have been without
|
|
// streaming at all.
|
|
if final == nil {
|
|
t.Fatal("no final run event; a client reading only the last event would have nothing")
|
|
}
|
|
if final["termination"] != "Completed" {
|
|
t.Errorf("final termination = %v", final["termination"])
|
|
}
|
|
if final["output"] != "Twelve events, mostly logins." {
|
|
t.Errorf("final output = %v", final["output"])
|
|
}
|
|
if final["runId"] == nil || final["runId"] == "" {
|
|
t.Error("the final event carries no run id")
|
|
}
|
|
}
|
|
|
|
func TestMiddlewareDoesNotSwallowFlush(t *testing.T) {
|
|
// The bug this pins cost an hour and produced no error anywhere.
|
|
//
|
|
// Two middlewares wrap the ResponseWriter to record a status and to
|
|
// intercept the mux's plain-text 404s. Both embed http.ResponseWriter,
|
|
// which inherits Write and WriteHeader and SILENTLY DROPS every optional
|
|
// interface underneath — Flusher among them. The SSE handler asked "can
|
|
// this flush?", was told no, and fell back to ordinary JSON: a correct,
|
|
// complete, entirely non-streaming response with nothing to indicate that
|
|
// streaming had been requested and quietly refused.
|
|
//
|
|
// Asserted through the whole middleware stack, because testing the handler
|
|
// alone is exactly what missed it.
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &streamingStub{pieces: []string{"a", "b"}}, nil)
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "test-agent")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@flush.test", "admin")
|
|
|
|
req, _ := http.NewRequest("POST", "/api/v1/agents/test-agent/runs",
|
|
strings.NewReader(`{"input":"hi"}`))
|
|
req.Header.Set("Content-Type", "application/json")
|
|
req.Header.Set("Accept", "text/event-stream")
|
|
req.AddCookie(admin.cookie)
|
|
|
|
rec := httptest.NewRecorder()
|
|
handler.ServeHTTP(rec, req)
|
|
|
|
if ct := rec.Header().Get("Content-Type"); !strings.Contains(ct, "text/event-stream") {
|
|
t.Fatalf("Content-Type is %q — a wrapper dropped Flusher and the stream fell back to JSON", ct)
|
|
}
|
|
if b := rec.Header().Get("X-Accel-Buffering"); b != "no" {
|
|
t.Errorf("X-Accel-Buffering is %q; a buffering proxy will hold the whole stream", b)
|
|
}
|
|
}
|
|
|
|
func TestAnOrdinaryRequestIsStillNotStreamed(t *testing.T) {
|
|
// Accept decides. A client that did not ask for a stream must not get one —
|
|
// it would be reading SSE frames as if they were a JSON body.
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &streamingStub{pieces: []string{"x"}}, nil)
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "test-agent")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@plain.test", "admin")
|
|
|
|
code, body := postRun(t, handler, admin, "test-agent", `{"input":"hi"}`)
|
|
if code != http.StatusOK {
|
|
t.Fatalf("status %d", code)
|
|
}
|
|
if body["termination"] != "Completed" || body["output"] != "x" {
|
|
t.Errorf("a plain request did not get a plain answer: %v", body)
|
|
}
|
|
}
|
|
|
|
func TestARequestedVersionReachesTheRuntime(t *testing.T) {
|
|
// §3's pin, at the seam. The frontend sends back the version its first
|
|
// answer carried; this asserts the field survives the request rather than
|
|
// being quietly dropped — which would look identical from outside until
|
|
// somebody published an edit mid-conversation.
|
|
h := testutil.New(t)
|
|
srv := runServer(t, h, &stubGateway{text: "answered"}, nil)
|
|
handler := srv.Handler()
|
|
|
|
orgID, adminID := seedOrgAdmin(t, h)
|
|
publishAgent(t, h.Pool, orgID, adminID, "test-agent")
|
|
admin := signInAs(t, handler, h.Pool, orgID, "admin", "admin@pin.test", "admin")
|
|
|
|
// Version 9 has no snapshot, so the run falls back to the current
|
|
// definition and says so — which is the observable proof the number
|
|
// travelled: an ignored field would produce no note at all.
|
|
code, body := postRun(t, handler, admin, "test-agent",
|
|
`{"input":"hello","agentVersion":9}`)
|
|
if code != http.StatusOK {
|
|
t.Fatalf("status %d: %v", code, body)
|
|
}
|
|
if body["termination"] != "Completed" {
|
|
t.Fatalf("termination = %v", body["termination"])
|
|
}
|
|
|
|
var runID, _ = body["runId"].(string)
|
|
req, _ := http.NewRequest("GET", "/api/v1/runs/"+runID, nil)
|
|
req.AddCookie(admin.cookie)
|
|
_, traj := doJSON(t, handler, req)
|
|
|
|
encoded, _ := json.Marshal(traj)
|
|
if !strings.Contains(string(encoded), "version_unavailable") {
|
|
t.Errorf("a pinned version with no snapshot left no trace in the trajectory; "+
|
|
"the field may have been dropped: %s", truncate(string(encoded), 400))
|
|
}
|
|
}
|
|
|
|
func truncate(s string, n int) string {
|
|
if len(s) <= n {
|
|
return s
|
|
}
|
|
return s[:n] + "…"
|
|
}
|