Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
42 commits
Select commit Hold shift + click to select a range
885fd90
feat(agent-vault): ship session activity from the proxy
saifsmailbox98 Sep 16, 2026
e103f77
fix(agent-vault): grow a session's activity buffer as records arrive
saifsmailbox98 Sep 23, 2026
90a6b73
fix(agent-vault): hold the proxy-wide activity byte cap across every …
saifsmailbox98 Sep 23, 2026
b3bb99b
fix(agent-vault): log the activity limit as an error, in the product'…
saifsmailbox98 Sep 23, 2026
6ad8783
fix(agent-vault): cap the method and port an agent can put in its own…
saifsmailbox98 Sep 23, 2026
6b3c563
fix(agent-vault): split activity chunks by size and count refused one…
saifsmailbox98 Sep 23, 2026
15b1cf8
fix(agent-vault): record an opaque request target as a blocked request
saifsmailbox98 Sep 23, 2026
62050d5
fix(agent-vault): upload activity only to https links, without redire…
saifsmailbox98 Sep 23, 2026
92ba968
fix(agent-vault): keep agent-chosen method and port out of proxy log …
saifsmailbox98 Sep 23, 2026
5f92a4e
fix(agent-vault): leave a posted chunk's drop count to its row
saifsmailbox98 Sep 23, 2026
b96369d
fix(agent-vault): send activity uploads as create-only and accept one…
saifsmailbox98 Sep 24, 2026
15ba4ae
chore(agent-vault): drop comments that restate the code
saifsmailbox98 Sep 24, 2026
93db017
chore(agent-vault): drop an unused chunk ID parser
saifsmailbox98 Sep 24, 2026
8bf2cf5
fix(agent-vault): warn with a count when a session's held activity is…
saifsmailbox98 Sep 24, 2026
21302ec
fix(agent-vault): say plainly when activity is refused for a wrong clock
saifsmailbox98 Sep 24, 2026
189e75c
fix(agent-vault): resume recording as soon as logging is back on
saifsmailbox98 Sep 24, 2026
fd8b61b
fix(agent-vault): give the final activity flush its own shutdown budget
saifsmailbox98 Sep 24, 2026
f86f38b
feat(agent-vault): report activity uploads on the heartbeat
saifsmailbox98 Sep 24, 2026
a08badf
fix(agent-vault): keep the final activity flush inside its shutdown b…
saifsmailbox98 Sep 24, 2026
9fde732
fix(agent-vault): flush each session's activity every minute, not eve…
saifsmailbox98 Sep 24, 2026
dfeb2f0
refactor(agent-vault): stop reporting activity uploads on the heartbeat
saifsmailbox98 Sep 25, 2026
9e5695f
feat(agent-vault): send each activity chunk's SHA-256 so the browser …
saifsmailbox98 Sep 25, 2026
5731392
improvement(agent-vault): seal activity chunks under sessionId and ch…
saifsmailbox98 Sep 25, 2026
46083b1
refactor(agent-vault): rename activity logs to session logs, on the w…
saifsmailbox98 Sep 26, 2026
96f740d
refactor(agent-vault): mint session log chunk ids as UUIDv7 instead o…
saifsmailbox98 Sep 26, 2026
653bb8e
test(agent-vault): cover held-byte release and asking for the session…
saifsmailbox98 Sep 27, 2026
fa698de
test(agent-vault): pin retry vs drop, the record limit per chunk, shu…
saifsmailbox98 Sep 27, 2026
5bb2f12
test(agent-vault): name the pinned session log vector after the brows…
saifsmailbox98 Sep 27, 2026
ff60021
refactor(agent-vault): drop the unread chunkId from the chunk reply a…
saifsmailbox98 Sep 27, 2026
cfa9f85
feat(agent-vault): send each chunk's SHA-256 on upload so S3 can chec…
saifsmailbox98 Sep 28, 2026
df3bfa7
fix(agent-vault): free shipped session log chunks, span each chunk fr…
saifsmailbox98 Sep 28, 2026
047d920
improvement(agent-vault): name S3's error code when the bucket refuse…
saifsmailbox98 Sep 28, 2026
8b1c1f5
fix(agent-vault): treat only Infisical's NotFound as a gone session, …
saifsmailbox98 Sep 28, 2026
46bb280
fix(agent-vault): retry a session log chunk on a 404 that isn't Infis…
saifsmailbox98 Sep 28, 2026
4fcb440
refactor(agent-vault): make the session log recorder easier to follow
saifsmailbox98 Sep 28, 2026
d282da4
refactor(agent-vault): order the session log recorder by stage and pu…
saifsmailbox98 Sep 28, 2026
0e62b8f
refactor(agent-vault): name the session log recorder's decisions and …
saifsmailbox98 Sep 28, 2026
37633ba
fix(agent-vault): keep forgotten session log spools in order so a ful…
saifsmailbox98 Sep 28, 2026
cd8e8bd
fix(agent-vault): let stop win over a due session log pass, ship sess…
saifsmailbox98 Sep 28, 2026
93d56dd
fix(agent-vault): seal session log chunks a round at a time, so a hea…
saifsmailbox98 Sep 29, 2026
0c4691d
fix(agent-vault): seal only what a session log pass started with, so …
saifsmailbox98 Sep 29, 2026
810aec6
fix(agent-vault): say session logs are disabled, without naming a pro…
saifsmailbox98 Sep 29, 2026
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
78 changes: 54 additions & 24 deletions packages/agentvault/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,11 +65,12 @@ type resolvedService struct {
}

type sessionEntry struct {
sessionID string
expiresAt *time.Time
services []*resolvedService
lastSeen time.Time
fetchedAt time.Time
sessionID string
expiresAt *time.Time
services []*resolvedService
sessionLog *sessionLogGrant
lastSeen time.Time
fetchedAt time.Time
}

// The map key is the sha256 of the token, never the token itself, so a heap dump yields no live credential.
Expand All @@ -79,7 +80,7 @@ func sessionKey(token string) string {
}

type sessionResolver interface {
resolve(sessionToken string) (*resolveResult, error)
resolve(sessionToken string, held *sessionLogGrant) (*resolveResult, error)
}

type sessionCache struct {
Expand Down Expand Up @@ -137,22 +138,41 @@ func isProxyTokenRejected(err error) bool {
return errors.As(err, &apiErr) && apiErr.Name == proxyTokenRejectedName
}

// The names on Infisical's own 404s and 401s. A route miss, a middlebox or an auth proxy in front of Infisical
// answers under another name or none.
const (
infisicalNotFoundName = "NotFound"
infisicalUnauthorizedName = "UnauthorizedError"
)

// Resolve answers 200, 401 or 404 by contract, and 401 with a name when the proxy's own token is the
// problem. Only those two statuses are a verdict on the session; anything else from a 4xx is a proxy-side
// fault or a middlebox, and the heartbeat classifier reads a 4xx the same way, so the two agree. A
// rejected proxy token is not a verdict on the session and is reported separately.
// problem. Only Infisical's own 401 (UnauthorizedError) or 404 (NotFound) is a verdict on the session; any
// other 4xx, an unnamed one included, is a proxy-side fault or a middlebox, and the heartbeat classifier reads a
// 4xx the same way, so the two agree. A rejected proxy token is not a verdict on the session and is reported
// separately.
func isSessionGone(err error) bool {
var apiErr *api.APIError
if errors.As(err, &apiErr) {
if isProxyTokenRejected(err) {
return false
}
return apiErr.StatusCode == http.StatusUnauthorized || apiErr.StatusCode == http.StatusNotFound
return (apiErr.StatusCode == http.StatusUnauthorized && apiErr.Name == infisicalUnauthorizedName) ||
(apiErr.StatusCode == http.StatusNotFound && apiErr.Name == infisicalNotFoundName)
}
return errors.Is(err, errSessionGone)
}

func (c *sessionCache) get(sessionToken string) ([]*resolvedService, error) {
services, _, err := c.lookup(sessionToken)
return services, err
}

type cacheLookup struct {
services []*resolvedService
sessionLog *sessionLogGrant
}

func (c *sessionCache) lookup(sessionToken string) ([]*resolvedService, *sessionLogGrant, error) {
key := sessionKey(sessionToken)

c.mu.Lock()
Expand All @@ -162,30 +182,30 @@ func (c *sessionCache) get(sessionToken string) ([]*resolvedService, error) {
delete(c.entries, key)
delete(c.tokens, key)
c.mu.Unlock()
return nil, errSessionGone
return nil, nil, errSessionGone
}
// Past the grace window the entry is a miss, so a stalled refresh loop cannot keep an old credential alive.
if time.Since(entry.fetchedAt) > c.grace() {
delete(c.entries, key)
delete(c.tokens, key)
} else {
entry.lastSeen = time.Now()
svcs := entry.services
svcs, grant := entry.services, entry.sessionLog
c.mu.Unlock()
return svcs, nil
return svcs, grant, nil
}
}
if refused, ok := c.refused[key]; ok {
if time.Now().Before(refused.until) {
c.mu.Unlock()
return nil, refused.err
return nil, nil, refused.err
}
delete(c.refused, key)
}
c.mu.Unlock()

resolved, err, _ := c.inflight.Do(key, func() (any, error) {
result, err := c.resolver.resolve(sessionToken)
result, err := c.resolver.resolve(sessionToken, nil)
if err != nil {
// A rejected proxy token is remembered too: the poll loop exits after two such heartbeats, but
// until then every agent request would otherwise cost a resolve.
Expand All @@ -201,19 +221,21 @@ func (c *sessionCache) get(sessionToken string) ([]*resolvedService, error) {
defer c.mu.Unlock()
c.evictIfFullLocked()
c.entries[key] = &sessionEntry{
sessionID: result.SessionID,
expiresAt: result.ExpiresAt,
services: result.Services,
lastSeen: time.Now(),
fetchedAt: time.Now(),
sessionID: result.SessionID,
expiresAt: result.ExpiresAt,
services: result.Services,
sessionLog: result.SessionLog,
lastSeen: time.Now(),
fetchedAt: time.Now(),
}
c.tokens[key] = sessionToken
return result.Services, nil
return cacheLookup{services: result.Services, sessionLog: result.SessionLog}, nil
})
if err != nil {
return nil, err
return nil, nil, err
}
return resolved.([]*resolvedService), nil
out := resolved.(cacheLookup)
return out.services, out.sessionLog, nil
}

func (c *sessionCache) evictIfFullLocked() {
Expand Down Expand Up @@ -289,7 +311,14 @@ func (c *sessionCache) refresh() {
}

func (c *sessionCache) refreshOne(key, token string) {
result, err := c.resolver.resolve(token)
c.mu.Lock()
var held *sessionLogGrant
if entry, ok := c.entries[key]; ok {
held = entry.sessionLog
}
c.mu.Unlock()

result, err := c.resolver.resolve(token, held)
if err != nil {
c.handleRefreshFailure(key, err)
return
Expand All @@ -301,6 +330,7 @@ func (c *sessionCache) refreshOne(key, token string) {
entry.sessionID = result.SessionID
entry.expiresAt = result.ExpiresAt
entry.services = result.Services
entry.sessionLog = result.SessionLog
entry.fetchedAt = time.Now()
}
}
Expand Down
38 changes: 22 additions & 16 deletions packages/agentvault/cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,16 +11,18 @@ import (
)

type stubResolver struct {
mu sync.Mutex
calls int
result *resolveResult
err error
delay time.Duration
mu sync.Mutex
calls int
result *resolveResult
err error
delay time.Duration
lastHeld *sessionLogGrant
}

func (s *stubResolver) resolve(string) (*resolveResult, error) {
func (s *stubResolver) resolve(_ string, held *sessionLogGrant) (*resolveResult, error) {
s.mu.Lock()
s.calls++
s.lastHeld = held
result, err, delay := s.result, s.err, s.delay
s.mu.Unlock()
time.Sleep(delay)
Expand Down Expand Up @@ -95,7 +97,7 @@ func TestRefreshDropsAGoneSessionImmediately(t *testing.T) {
t.Fatalf("get: %v", err)
}

resolver.err = &api.APIError{StatusCode: status}
resolver.err = &api.APIError{StatusCode: status, Name: map[int]string{401: infisicalUnauthorizedName, 404: infisicalNotFoundName}[status]}
cache.refresh()

if len(cache.entries) != 0 {
Expand Down Expand Up @@ -196,27 +198,31 @@ func TestStaleEntryIsNotServedFromTheCache(t *testing.T) {
}
}

// Only the statuses resolve answers by contract end a session. A 400, 403 or 422 cannot come from resolve
// itself, so it is a proxy-side fault or a middlebox and rides the grace window like an outage.
// Only the refusals resolve answers by contract end a session. A 400, 403 or 422, or a 404 that is not
// Infisical's NotFound, cannot come from resolve itself, so it is a proxy-side fault or a middlebox and rides
// the grace window like an outage.
func TestRefreshTreatsOnlyTheContractRefusalsAsTerminal(t *testing.T) {
for _, tc := range []struct {
status int
name string
kept bool
}{
{401, false}, {404, false},
{400, true}, {403, true}, {405, true}, {407, true}, {422, true},
{408, true}, {429, true}, {500, true}, {502, true},
{401, infisicalUnauthorizedName, false}, {404, infisicalNotFoundName, false},
{401, "", true}, {401, "Unauthorized", true},
{404, "", true}, {404, "Not Found", true},
{400, "", true}, {403, "", true}, {405, "", true}, {407, "", true}, {422, "", true},
{408, "", true}, {429, "", true}, {500, "", true}, {502, "", true},
} {
resolver := &stubResolver{result: &resolveResult{SessionID: "s1", Services: []*resolvedService{serviceWithSecret("v")}}}
cache := newTestCache(resolver)
if _, err := cache.get("tok"); err != nil {
t.Fatal(err)
}
resolver.err = &api.APIError{StatusCode: tc.status}
resolver.err = &api.APIError{StatusCode: tc.status, Name: tc.name}
cache.refresh()
_, kept := cache.get("tok")
if (kept == nil) != tc.kept {
t.Fatalf("status %d: credential still served = %v, want %v", tc.status, kept == nil, tc.kept)
t.Fatalf("status %d %q: credential still served = %v, want %v", tc.status, tc.name, kept == nil, tc.kept)
}
}
}
Expand Down Expand Up @@ -305,7 +311,7 @@ func TestOnlyOneRefreshRunsAtATime(t *testing.T) {
}

func TestADefinitiveRefusalIsNotReResolvedEveryRequest(t *testing.T) {
resolver := &stubResolver{err: &api.APIError{StatusCode: 404, Name: "NotFound"}}
resolver := &stubResolver{err: &api.APIError{StatusCode: 404, Name: infisicalNotFoundName}}
cache := newTestCache(resolver)

for i := 0; i < 20; i++ {
Expand Down Expand Up @@ -343,7 +349,7 @@ func TestAnOutageIsNotRememberedAsARefusal(t *testing.T) {
}

func TestARefusalExpires(t *testing.T) {
resolver := &stubResolver{err: &api.APIError{StatusCode: 404}}
resolver := &stubResolver{err: &api.APIError{StatusCode: 404, Name: infisicalNotFoundName}}
cache := newTestCache(resolver)
_, _ = cache.get("agv_dead")
cache.mu.Lock()
Expand Down
16 changes: 10 additions & 6 deletions packages/agentvault/policy.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ var errBodyUnreadable = errors.New("could not read the request body")

func checkServicePolicy(svc *resolvedService, req *http.Request) error {
if !svc.allowsMethod(req.Method) {
return fmt.Errorf("service %q does not allow %s: %w", svc.name, req.Method, errPolicyBlocked)
return fmt.Errorf("service %q does not allow %s: %w", svc.name, truncateLogged(req.Method, maxLoggedMethodLen), errPolicyBlocked)
}
if len(svc.allowedPathPrefixes) > 0 {
path := requestPath(req)
Expand Down Expand Up @@ -103,8 +103,8 @@ func escapeInvalidPathBytes(raw string) string {
func requestPath(req *http.Request) string {
path := req.URL.EscapedPath()
if path == "" {
// forwardHTTP refuses an opaque target before this runs, so the branch is a floor under that check
// rather than a shape expected here. A genuinely empty path is the root.
// forward refuses an opaque target before any policy reads the path, so this branch is what the log
// line and the session log record show for one. A genuinely empty path is the root.
if req.URL.Opaque != "" {
return req.URL.Opaque
}
Expand All @@ -114,10 +114,14 @@ func requestPath(req *http.Request) string {
}

func truncatePath(path string) string {
if len(path) > maxLoggedPathLen {
return path[:maxLoggedPathLen] + "...[truncated]"
return truncateLogged(path, maxLoggedPathLen)
}

func truncateLogged(value string, limit int) string {
if len(value) > limit {
return value[:limit] + "...[truncated]"
}
return path
return value
}

// Never decodes: anything whose meaning depends on the upstream's normalisation is refused outright, so the
Expand Down
11 changes: 11 additions & 0 deletions packages/agentvault/policy_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,17 @@ func TestMethodPolicy(t *testing.T) {
}
})

t.Run("a refused method is capped before it reaches the log", func(t *testing.T) {
svc := serviceWithPolicy([]string{"GET"}, nil)
err := checkServicePolicy(svc, requestTo(t, strings.Repeat("A", 1<<20), "/x"))
if !errors.Is(err, errPolicyBlocked) {
t.Fatalf("an unlisted method should be blocked, got %v", err)
}
if len(err.Error()) > 200 {
t.Fatalf("the refusal carries %d bytes; the agent's method would fill the proxy log", len(err.Error()))
}
})

t.Run("a lower-case method is folded rather than blocked", func(t *testing.T) {
svc := serviceWithPolicy([]string{"GET"}, nil)
req := requestTo(t, "GET", "/x")
Expand Down
Loading
Loading