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
35 changes: 19 additions & 16 deletions internal/gateway/archive.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
package gateway

// The request archive: with it on (settings' RequestArchive), each call the
// gateway serves is kept — the headers and bodies both ways, every secret
// taken out first — in the user's own S3 bucket, the one sync keeps its
// backup in, at <prefix>/magpie/archive/<date>/<id>.json: the call's date,
// in UTC, and an id of its own, sent back on the response as
// gateway serves is kept — the headers and bodies both ways, with every
// secret, and whatever the user's own masking rules, words and personal
// data cover, taken out first — in the user's own S3 bucket, the one sync
// keeps its backup in, at <prefix>/magpie/archive/<date>/<id>.json: the
// call's date, in UTC, and an id of its own, sent back on the response as
// X-Magpie-Archive-Id. The call, and its rows in the usage log, carry
// "<date>/<id>", for the Gateway and Usage pages to read it back by.
//
Expand Down Expand Up @@ -74,6 +75,7 @@ type wire struct {
reqFull *spool
res *captureResponseWriter
to Putter
o redact.Options // what its secrets are taken out with
}

// archiving is r's wire when the archive is on and has a bucket, nil
Expand All @@ -93,7 +95,7 @@ func archiving(r *http.Request, res *captureResponseWriter, start time.Time, bod
res.Header().Set(ArchiveHeader, id)
res.full = &spool{limit: archiveLimit()}
return &wire{name: start.UTC().Format("2006-01-02") + "/" + id, method: r.Method, path: r.URL.Path,
query: r.URL.RawQuery, req: r.Header.Clone(), body: body, res: res, to: to}
query: r.URL.RawQuery, req: r.Header.Clone(), body: body, res: res, to: to, o: scrubOptions()}
}

// archiveName is where the archive keeps c, "<date>/<id>", or "" when it
Expand Down Expand Up @@ -210,22 +212,23 @@ func upload(j archiveJob) {
c := j.c
c.wire = nil
c.RequestBody, c.ResponseBody = "", ""
o := j.w.o
reqBody, err1 := j.req.read()
resBody, err2 := j.res.read()
if err := errors.Join(err1, err2); err != nil {
archiveFailed(j.w.name, err)
return
}
c.Error, c.Fallback = redact.Scrub(c.Error), redact.Scrub(c.Fallback)
c.Error, c.Fallback = redact.ScrubWith(c.Error, o), redact.ScrubWith(c.Fallback, o)
path := j.w.path
if j.w.query != "" {
path += "?" + scrubQuery(j.w.query)
path += "?" + scrubQuery(j.w.query, o)
}
a := Archived{ID: j.w.name, Call: c,
Request: ArchivePart{Method: j.w.method, Path: path, Headers: scrubHeaders(j.w.req),
Body: string(redact.ScrubJSON(reqBody)), Size: j.req.size, Truncated: j.req.cut()},
Response: ArchivePart{Status: c.Status, Headers: scrubHeaders(j.resHead),
Body: string(redact.ScrubJSON(resBody)), Size: j.res.size, Truncated: j.res.cut()},
Request: ArchivePart{Method: j.w.method, Path: path, Headers: scrubHeaders(j.w.req, o),
Body: string(redact.ScrubJSONWith(reqBody, o)), Size: j.req.size, Truncated: j.req.cut()},
Response: ArchivePart{Status: c.Status, Headers: scrubHeaders(j.resHead, o),
Body: string(redact.ScrubJSONWith(resBody, o)), Size: j.res.size, Truncated: j.res.cut()},
}
// as it reads in the bucket: a query's & and a body's <tags> as they are
var buf bytes.Buffer
Expand Down Expand Up @@ -263,12 +266,12 @@ func ArchiveError() (string, time.Time) {
return archiveErr, archiveErrAt
}

func scrubHeaders(h http.Header) map[string][]string {
func scrubHeaders(h http.Header, o redact.Options) map[string][]string {
out := make(map[string][]string, len(h))
for k, vs := range h {
s := make([]string, len(vs))
for i, v := range vs {
s[i] = redact.ScrubHeader(k, v)
s[i] = redact.ScrubHeaderWith(k, v, o)
}
out[k] = s
}
Expand All @@ -277,17 +280,17 @@ func scrubHeaders(h http.Header) map[string][]string {

// scrubQuery is a query string with the values of secret-named fields
// (key=, access_token=) and any secret in the others taken out.
func scrubQuery(q string) string {
func scrubQuery(q string, o redact.Options) string {
vs, err := url.ParseQuery(q)
if err != nil {
return redact.Scrub(q)
return redact.ScrubWith(q, o)
}
for k, v := range vs {
for i := range v {
if redact.SecretName(k) {
v[i] = redact.Scrubbed
} else {
v[i] = redact.Scrub(v[i])
v[i] = redact.ScrubWith(v[i], o)
}
}
}
Expand Down
87 changes: 87 additions & 0 deletions internal/gateway/archive_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import (
"sync"
"testing"

"github.com/yetone/magpie/internal/redact"
"github.com/yetone/magpie/internal/settings"
)

Expand Down Expand Up @@ -276,3 +277,89 @@ func keys(m map[string][]byte) []string {
}
return out
}

// relayKey is a relay's key of a format of the user's own: magpie's rules
// know nothing of it, so only their masking rule finds it (#195).
const relayKey = "rz_RelayKey1234567"

// relayVendor answers with the relay key in its reply, as a vendor that
// names a key back does.
type relayVendor struct{}

func (relayVendor) ServeHTTP(w http.ResponseWriter, r *http.Request) {
req, _ := io.ReadAll(r.Body)
key := relayKey
// echo the key the agent sent, so a check using one of its own reads it
// back: masking off, the vendor saw the user's rule take nothing out
if i := strings.Index(string(req), "rz_"); i >= 0 {
rest := string(req)[i:]
if j := strings.IndexAny(rest, `"`); j > 0 {
key = rest[:j]
}
}
w.Header().Set("Content-Type", "application/json")
io.WriteString(w, `{"id":"c1","object":"chat.completion","model":"m1","choices":[{"index":0,"message":{"role":"assistant","content":"using `+key+`"},"finish_reason":"stop"}],"usage":{"prompt_tokens":3,"completion_tokens":1}}`)
}

// Every secret of magpie's own goes out of the archive whether masking is
// on for the vendor or not: what the archive keeps is sent nowhere, so a
// magpie rule's match goes from it either way. A masking rule of the user's
// own is left behind: mask gates it on Mask secrets, so forcing the secrets
// on for what is kept would switch their rule on where the vendor side left
// it off. With masking on the vendor was sent a placeholder and the archive
// keeps that one; with it off the agent's own text is what the archive
// keeps, the relay key and all, since no rule of magpie's knows rz_.
func TestArchiveKeepsTheUsersRules(t *testing.T) {
for _, tc := range []struct {
name string
redact bool
key string
}{{"masking_on", true, "rz_RelayKey1234567"}, {"masking_off", false, "rz_RelayKey7654321"}} {
t.Run(tc.name, func(t *testing.T) {
fresh(t)
serveOn(t, "fake", "k", []string{"m1"}, relayVendor{})
if err := settings.Save(settings.Settings{RequestArchive: true, Redact: tc.redact,
RedactRules: []redact.Rule{{Kind: "RELAY", Prefix: "rz_"}}}); err != nil {
t.Fatal(err)
}
b := &memBucket{objs: map[string][]byte{}}
archiveTo(t, b)
s := New()
body := `{"model":"fake/m1","messages":[{"role":"user","content":"send it to ` + tc.key + `"}]}`
rec := httptest.NewRecorder()
s.Handler().ServeHTTP(rec, httptest.NewRequest("POST", "/v1/chat/completions", strings.NewReader(body)))
archivePending.Wait()
if rec.Code != 200 || !strings.Contains(rec.Body.String(), tc.key) {
t.Fatalf("%d %s", rec.Code, rec.Body)
}
data, ok := b.objs["archive/"+s.Recent()[0].Archive+".json"]
if !ok || len(b.objs) != 1 {
t.Fatalf("uploaded %v", keys(b.objs))
}
var a Archived
if err := json.Unmarshal(data, &a); err != nil {
t.Fatal(err)
}
// masking on, the vendor was sent a placeholder and the archive
// keeps that one; off, the user's rule took nothing out of what
// the vendor was sent, and no rule of magpie's knows rz_, so the
// key is in the archived request as it was in the agent's
want := tc.key
if tc.redact {
want = "{{RELAY_"
}
if !strings.Contains(a.Request.Body, want) {
t.Errorf("want %s in the archived request:\n%s", want, a.Request.Body)
}
if tc.redact == strings.Contains(a.Request.Body, tc.key) {
t.Errorf("masking %v: the relay key in the archived request:\n%s", tc.redact, a.Request.Body)
}
// the reply was masked on its way back only where the vendor
// echoed the placeholder: with masking on the key is restored
// for the agent, and with it off the vendor never saw one
if !tc.redact && !strings.Contains(a.Response.Body, tc.key) {
t.Errorf("the relay key is not in the archived reply, and no rule knows it:\n%s", a.Response.Body)
}
})
}
}
2 changes: 1 addition & 1 deletion internal/gateway/capture.go
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ func bodyForExport(body string, cut bool) string {
if body == "" {
return ""
}
body = string(redact.ScrubJSON([]byte(body)))
body = string(redact.ScrubJSONWith([]byte(body), scrubOptions()))
if cut {
body += usage.BodyCut
}
Expand Down
33 changes: 33 additions & 0 deletions internal/gateway/capture_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,9 @@ import (
"net/http/httptest"
"strings"
"testing"

"github.com/yetone/magpie/internal/redact"
"github.com/yetone/magpie/internal/settings"
)

func TestCaptureRequestBodyLimit(t *testing.T) {
Expand All @@ -25,3 +28,33 @@ func TestCaptureResponseWriterPreservesResponse(t *testing.T) {
t.Fatalf("response = %q", got)
}
}

// The bodies the OTLP export carries lose every secret of magpie's own
// whether masking is on for the vendor or not, since what is sent to the
// collector is kept nowhere (#195). A masking rule of the user's own is
// left behind: mask gates it on Mask secrets, so forcing the secrets on
// here would switch their rule on where the vendor side left it off.
func TestBodyForExportKeepsTheUsersRules(t *testing.T) {
fresh(t)
if err := settings.Save(settings.Settings{Redact: true,
RedactRules: []redact.Rule{{Kind: "RELAY", Prefix: "rz_"}}}); err != nil {
t.Fatal(err)
}
// a key of this test's own: another test's, or another case's, may
// already have masked its own, and a value masked once is masked again
// wherever a request has it (known.go), whatever the options say
const key = "rz_ExportKey1234567"
body := `{"choices":[{"message":{"content":"using ` + key + `"}}]}`
if got := bodyForExport(body, false); !strings.Contains(got, key) {
t.Errorf("exported body: %s", got)
}
// masking off, the user's rule is off with it, as the vendor side has it,
// and no rule of magpie's knows rz_, so the key is what is kept
if err := settings.Save(settings.Settings{
RedactRules: []redact.Rule{{Kind: "RELAY", Prefix: "rz_"}}}); err != nil {
t.Fatal(err)
}
if got := bodyForExport(body, false); !strings.Contains(got, key) {
t.Errorf("exported body with masking off took the rule's match out: %s", got)
}
}
22 changes: 18 additions & 4 deletions internal/gateway/redact.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ import (
// writes what the wrapper still holds. Nothing masked, nothing wrapped.
func redacted(w http.ResponseWriter, body []byte) (http.ResponseWriter, []byte, func()) {
o := redactionOptions()
if !o.Secrets && !o.Personal && len(o.Words) == 0 {
if !o.Secrets && !o.Personal && len(o.Words) == 0 && len(o.Rules) == 0 {
return w, body, func() {}
}
masked, n := redact.MaskJSON(body, o)
Expand All @@ -35,9 +35,23 @@ func redactedPrompt(w http.ResponseWriter, prompt string) (http.ResponseWriter,
return rw, masked, rw.Finish
}

func redactionOptions() redact.Options {
st := settings.Load()
return redact.Options{Secrets: st.Redact, Personal: st.RedactPersonal, Kinds: st.RedactKinds, Words: st.RedactWords, Rules: st.RedactRules}
func redactionOptions() redact.Options { return settings.Load().Redaction() }

// scrubOptions is what the request archive and the OTLP bodies take out of
// what they keep: what the user asked masked, and every secret beside it.
// What they keep is sent nowhere — the user's own bucket, or their
// collector — so a secret goes from it whether or not Mask secrets is on,
// which is the archive's contract. Secrets forced this way are magpie's
// own rules, and nothing else: a masking rule of the user's own stays
// gated on Mask secrets as settings documents it, since forcing Secrets on
// would switch those on with them (redact.mask gates the user's rules on
// Secrets). Their words and personal data follow their own switches as
// before.
func scrubOptions() redact.Options {
o := redactionOptions()
// the user's rules are left out: they are gated on Mask secrets, which
// is settings' documented meaning of them
return redact.Options{Secrets: true, Personal: o.Personal, Words: o.Words}
}

// unredactedRoute says a request resolved to p, or to the group whose
Expand Down
32 changes: 23 additions & 9 deletions internal/redact/scrub.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,52 +27,66 @@ func SecretName(name string) bool {
}

// Scrub takes the secrets Mask finds out of s.
func Scrub(s string) string {
out, _ := mask(s, Options{Secrets: true}, scrubbed)
func Scrub(s string) string { return ScrubWith(s, Options{Secrets: true}) }

// ScrubWith takes what o covers out of s, as Scrub takes the secrets: the
// user's own rules, words and personal data go with them, for what is kept
// rather than sent — the request archive, the OTLP bodies (#195).
func ScrubWith(s string, o Options) string {
out, _ := mask(s, o, scrubbed)
return out
}

// ScrubHeader is a header's value with its secrets taken out: all of it
// for a header named for one or a credential (Bearer, Basic), else what
// Scrub finds in it.
func ScrubHeader(name, value string) string {
return ScrubHeaderWith(name, value, Options{Secrets: true})
}

// ScrubHeaderWith is ScrubHeader with what o covers taken out of a value
// its name says nothing about.
func ScrubHeaderWith(name, value string, o Options) string {
v := strings.ToLower(strings.TrimSpace(value))
if SecretName(name) || strings.HasPrefix(v, "bearer ") || strings.HasPrefix(v, "basic ") {
return Scrubbed
}
return Scrub(value)
return ScrubWith(value, o)
}

// ScrubJSON is a body with its secrets taken out: a JSON one string by
// string, a field named for a secret wholly; a stream's events each as
// JSON, line by line; anything else as text. A picture's data is left as
// it is.
func ScrubJSON(body []byte) []byte {
if out, ok := scrubJSON(body); ok {
func ScrubJSON(body []byte) []byte { return ScrubJSONWith(body, Options{Secrets: true}) }

// ScrubJSONWith is ScrubJSON with what o covers taken out.
func ScrubJSONWith(body []byte, o Options) []byte {
if out, ok := scrubJSON(body, o); ok {
return out
}
lines := strings.SplitAfter(string(body), "\n")
for i, l := range lines {
if rest, ok := strings.CutPrefix(l, "data:"); ok {
trimmed := strings.TrimSpace(rest)
if out, ok := scrubJSON([]byte(trimmed)); ok && trimmed != "" {
if out, ok := scrubJSON([]byte(trimmed), o); ok && trimmed != "" {
lines[i] = "data: " + string(out) + l[len(strings.TrimRight(l, "\r\n")):]
continue
}
}
lines[i] = Scrub(l)
lines[i] = ScrubWith(l, o)
}
return []byte(strings.Join(lines, ""))
}

func scrubJSON(b []byte) ([]byte, bool) {
func scrubJSON(b []byte, o Options) ([]byte, bool) {
return walk(b, func(_, key, s string) string {
switch {
case strings.HasPrefix(s, "data:"):
return s
case key != "" && SecretName(key):
return Scrubbed
}
return Scrub(s)
return ScrubWith(s, o)
})
}
Loading
Loading