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
6 changes: 6 additions & 0 deletions .claude/e2e/exploit_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,21 @@ func TestPublicExploitSmoke(t *testing.T) {

root := findProjectRoot()
require.NotEmpty(t, root, "could not find project root")
kitchenURL := getEnvOrFile("KITCHEN_URL", e2eEnvPath)
require.NotEmpty(t, kitchenURL, "KITCHEN_URL required")
authToken := getEnvOrFile("AUTH_TOKEN", e2eEnvPath)
require.NotEmpty(t, authToken, "AUTH_TOKEN required")

require.NoError(t, resetE2EWorkspace())
require.NoError(t, waitForKitchenHealth(root, 90*time.Second))
graphProbes := connectGraphSmokeProbes(t, kitchenURL, authToken)
require.NoError(t, writeConfig(token, targetRepo))
require.NoError(t, restartCounter("e2e-smoke"))

tmux := newTmuxController(tmuxSessionName)

waitForReconPhase(t, tmux)
verifyGraphPublicationAfterAnalysis(t, graphProbes)
requireContent(t, tmux, "whooli", 15*time.Second, "Attack tree should show target org")
requireContent(t, tmux, "xyz", 15*time.Second, "Attack tree should show target repo")

Expand Down
224 changes: 224 additions & 0 deletions .claude/e2e/graph.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,224 @@
// Copyright (C) 2026 boostsecurity.io
// SPDX-License-Identifier: AGPL-3.0-or-later

//go:build e2e

package e2e

import (
"context"
"encoding/json"
"fmt"
"net/url"
"testing"
"time"

"github.com/coder/websocket"
"github.com/coder/websocket/wsjson"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

const graphProbeTimeout = 30 * time.Second

type graphWireMessage struct {
Type string `json:"type"`
Data json.RawMessage `json:"data"`
}

type graphWireSnapshot struct {
Revision uint64 `json:"revision"`
Mode string `json:"mode"`
TotalNodes int `json:"total_nodes"`
TotalEdges int `json:"total_edges"`
Nodes []json.RawMessage `json:"nodes"`
Edges []json.RawMessage `json:"edges"`
}

type graphWireDelta struct {
BaseRevision uint64 `json:"base_revision"`
Revision uint64 `json:"revision"`
}

type graphWireFence struct {
Revision uint64 `json:"revision"`
}

type graphProbe struct {
connection *websocket.Conn
mode string
revision uint64
}

type graphSmokeProbes struct {
full *graphProbe
filtered *graphProbe
auto *graphProbe
}

func connectGraphSmokeProbes(t *testing.T, kitchenURL, authToken string) *graphSmokeProbes {
t.Helper()

return &graphSmokeProbes{
full: connectGraphProbe(t, kitchenURL, authToken, "full"),
filtered: connectGraphProbe(t, kitchenURL, authToken, "filtered"),
auto: connectGraphProbe(t, kitchenURL, authToken, "auto"),
}
}

func connectGraphProbe(t *testing.T, kitchenURL, authToken, mode string) *graphProbe {
t.Helper()

endpoint, err := graphWebSocketURL(kitchenURL, authToken, mode)
require.NoError(t, err)
ctx, cancel := context.WithTimeout(context.Background(), graphProbeTimeout)
defer cancel()
connection, _, err := websocket.Dial(ctx, endpoint, nil)
require.NoError(t, err)
connection.SetReadLimit(-1)
t.Cleanup(func() { _ = connection.Close(websocket.StatusNormalClosure, "") })

probe := &graphProbe{connection: connection, mode: mode}
message := probe.read(t)
require.Equal(t, "snapshot", message.Type, "%s mode must start from a complete snapshot", mode)
snapshot := decodeGraphData[graphWireSnapshot](t, message)
require.Equal(t, uint64(0), snapshot.Revision, "%s mode should connect to the purged Pantry", mode)
require.Zero(t, snapshot.TotalNodes)
require.Zero(t, snapshot.TotalEdges)
probe.revision = snapshot.Revision
return probe
}

func graphWebSocketURL(kitchenURL, authToken, mode string) (string, error) {
endpoint, err := url.Parse(kitchenURL)
if err != nil {
return "", fmt.Errorf("parse Kitchen URL: %w", err)
}
switch endpoint.Scheme {
case "http":
endpoint.Scheme = "ws"
case "https":
endpoint.Scheme = "wss"
default:
return "", fmt.Errorf("unsupported Kitchen URL scheme %q", endpoint.Scheme)
}
endpoint.Path = "/graph/ws"
query := endpoint.Query()
query.Set("token", authToken)
query.Set("mode", mode)
endpoint.RawQuery = query.Encode()
return endpoint.String(), nil
}

func verifyGraphPublicationAfterAnalysis(t *testing.T, probes *graphSmokeProbes) {
t.Helper()

fullSnapshot := probes.full.refreshAfterFence(t)
require.Equal(t, "full", fullSnapshot.Mode)
require.NotZero(t, fullSnapshot.TotalNodes)
require.Len(t, fullSnapshot.Nodes, fullSnapshot.TotalNodes)
require.Len(t, fullSnapshot.Edges, fullSnapshot.TotalEdges)

filteredSnapshot := probes.filtered.waitForSnapshot(t, fullSnapshot.Revision)
require.Equal(t, "filtered", filteredSnapshot.Mode)
require.NotZero(t, filteredSnapshot.TotalNodes)
assert.LessOrEqual(t, len(filteredSnapshot.Nodes), filteredSnapshot.TotalNodes)
assert.LessOrEqual(t, len(filteredSnapshot.Edges), filteredSnapshot.TotalEdges)

autoSnapshot := probes.auto.waitForSnapshot(t, fullSnapshot.Revision)
require.NotZero(t, autoSnapshot.TotalNodes)
assert.Contains(t, []string{"full", "filtered"}, autoSnapshot.Mode)
assert.LessOrEqual(t, len(autoSnapshot.Nodes), autoSnapshot.TotalNodes)
assert.LessOrEqual(t, len(autoSnapshot.Edges), autoSnapshot.TotalEdges)
}

func (p *graphProbe) refreshAfterFence(t *testing.T) graphWireSnapshot {
t.Helper()

minimumRevision := p.revision
for {
message := p.read(t)
switch message.Type {
case "delta":
delta := decodeGraphData[graphWireDelta](t, message)
require.Equal(t, p.revision, delta.BaseRevision, "full-mode deltas must be contiguous before a fence")
require.Greater(t, delta.Revision, p.revision)
p.revision = delta.Revision
case "snapshot_required":
fence := decodeGraphData[graphWireFence](t, message)
require.Greater(t, fence.Revision, p.revision)
minimumRevision = fence.Revision
p.requestSnapshot(t)
return p.waitForFullSnapshot(t, minimumRevision)
default:
t.Fatalf("full mode received unexpected graph message %q before committed-state fence", message.Type)
}
}
}

func (p *graphProbe) waitForFullSnapshot(t *testing.T, minimumRevision uint64) graphWireSnapshot {
t.Helper()

for {
message := p.read(t)
switch message.Type {
case "snapshot":
snapshot := decodeGraphData[graphWireSnapshot](t, message)
require.GreaterOrEqual(t, snapshot.Revision, minimumRevision)
p.revision = snapshot.Revision
return snapshot
case "delta":
delta := decodeGraphData[graphWireDelta](t, message)
minimumRevision = max(minimumRevision, delta.Revision)
case "snapshot_required":
fence := decodeGraphData[graphWireFence](t, message)
minimumRevision = max(minimumRevision, fence.Revision)
default:
t.Fatalf("full mode received unexpected graph message %q while refreshing", message.Type)
}
}
}

func (p *graphProbe) waitForSnapshot(t *testing.T, minimumRevision uint64) graphWireSnapshot {
t.Helper()

for {
message := p.read(t)
require.Equal(t, "snapshot", message.Type, "%s mode must receive complete snapshots while refreshing", p.mode)
snapshot := decodeGraphData[graphWireSnapshot](t, message)
require.GreaterOrEqual(t, snapshot.Revision, p.revision)
p.revision = snapshot.Revision
if snapshot.Revision >= minimumRevision {
return snapshot
}
}
}

func (p *graphProbe) requestSnapshot(t *testing.T) {
t.Helper()

ctx, cancel := context.WithTimeout(context.Background(), graphProbeTimeout)
defer cancel()
require.NoError(t, wsjson.Write(ctx, p.connection, map[string]any{
"type": "snapshot_request",
"data": map[string]string{"mode": p.mode},
}))
}

func (p *graphProbe) read(t *testing.T) graphWireMessage {
t.Helper()

ctx, cancel := context.WithTimeout(context.Background(), graphProbeTimeout)
defer cancel()
var message graphWireMessage
require.NoError(t, wsjson.Read(ctx, p.connection, &message), "%s graph stream did not publish in time", p.mode)
return message
}

func decodeGraphData[T any](t *testing.T, message graphWireMessage) T {
t.Helper()

var data T
require.NoError(t, json.Unmarshal(message.Data, &data))
return data
}
68 changes: 49 additions & 19 deletions internal/kitchen/browser_assets/graph.js
Original file line number Diff line number Diff line change
Expand Up @@ -31,12 +31,14 @@ SPDX-License-Identifier: AGPL-3.0-or-later

let cy;
let ws;
let graphVersion = 0;
let graphRevision = 0;
let reconnectAttempts = 0;
let requestedGraphMode = readGraphMode();
let resolvedGraphMode = 'full';
let graphMeta = { totalNodes: 0, totalEdges: 0, largeGraph: false, filterDescription: '' };
let filteredRefreshTimer = null;
let minimumSnapshotRevision = 0;
let snapshotRequestPending = false;
let awaitingInitialSnapshot = true;

function initCytoscape() {
cy = cytoscape({
Expand Down Expand Up @@ -152,6 +154,9 @@ SPDX-License-Identifier: AGPL-3.0-or-later
wsUrl += '?' + wsParams.toString();

updateConnectionStatus('connecting');
awaitingInitialSnapshot = true;
snapshotRequestPending = false;
minimumSnapshotRevision = Math.max(minimumSnapshotRevision, graphRevision);
ws = new WebSocket(wsUrl);

ws.onopen = function() {
Expand All @@ -165,6 +170,7 @@ SPDX-License-Identifier: AGPL-3.0-or-later
};

ws.onclose = function() {
snapshotRequestPending = false;
updateConnectionStatus('disconnected');
scheduleReconnect();
};
Expand All @@ -188,13 +194,24 @@ SPDX-License-Identifier: AGPL-3.0-or-later
case 'delta':
handleDelta(msg.data);
break;
case 'snapshot_required':
requireSnapshot(msg.data.revision);
break;
case 'pong':
break;
}
}

function handleSnapshot(data) {
graphVersion = data.version;
snapshotRequestPending = false;
if (!Number.isSafeInteger(data.revision) || data.revision < minimumSnapshotRevision || data.revision < graphRevision) {
requestRequiredSnapshot();
return;
}

graphRevision = data.revision;
minimumSnapshotRevision = graphRevision;
awaitingInitialSnapshot = false;
resolvedGraphMode = data.mode || 'full';
graphMeta = {
totalNodes: data.total_nodes || 0,
Expand Down Expand Up @@ -238,12 +255,21 @@ SPDX-License-Identifier: AGPL-3.0-or-later
}

function handleDelta(data) {
if (resolvedGraphMode !== 'full' || prefersFilteredSnapshots()) {
scheduleFilteredRefresh();
if (!Number.isSafeInteger(data.base_revision) || !Number.isSafeInteger(data.revision)) {
requireSnapshot(graphRevision + 1);
return;
}
if (data.revision <= graphRevision) {
return;
}
if (awaitingInitialSnapshot || requestedGraphMode !== 'full' || resolvedGraphMode !== 'full') {
requireSnapshot(data.revision);
return;
}
if (minimumSnapshotRevision > graphRevision || data.base_revision !== graphRevision) {
requireSnapshot(data.revision);
return;
}

graphVersion = data.version;

(data.added_nodes || []).forEach(node => {
cy.add({
Expand Down Expand Up @@ -297,6 +323,8 @@ SPDX-License-Identifier: AGPL-3.0-or-later
if (!graphMeta.largeGraph && (data.added_nodes || []).length > 2) {
runLayout();
}
graphRevision = data.revision;
minimumSnapshotRevision = graphRevision;
updateStats();
}

Expand Down Expand Up @@ -344,10 +372,10 @@ SPDX-License-Identifier: AGPL-3.0-or-later
const edges = cy.edges().length;
if (resolvedGraphMode === 'filtered') {
document.getElementById('stats').textContent =
'Filtered: ' + nodes + '/' + graphMeta.totalNodes + ' nodes, ' + edges + '/' + graphMeta.totalEdges + ' edges (v' + graphVersion + ')';
'Filtered: ' + nodes + '/' + graphMeta.totalNodes + ' nodes, ' + edges + '/' + graphMeta.totalEdges + ' edges (r' + graphRevision + ')';
return;
}
document.getElementById('stats').textContent = 'Full: ' + nodes + ' nodes, ' + edges + ' edges (v' + graphVersion + ')';
document.getElementById('stats').textContent = 'Full: ' + nodes + ' nodes, ' + edges + ' edges (r' + graphRevision + ')';
}

function updateConnectionStatus(status) {
Expand Down Expand Up @@ -387,26 +415,28 @@ SPDX-License-Identifier: AGPL-3.0-or-later

function requestSnapshot(mode) {
if (!ws || ws.readyState !== WebSocket.OPEN) {
return;
return false;
}
ws.send(JSON.stringify({
type: 'snapshot_request',
data: { mode: mode }
}));
snapshotRequestPending = true;
return true;
}

function scheduleFilteredRefresh() {
if (filteredRefreshTimer) {
return;
function requireSnapshot(revision) {
if (Number.isSafeInteger(revision)) {
minimumSnapshotRevision = Math.max(minimumSnapshotRevision, revision);
}
filteredRefreshTimer = window.setTimeout(function() {
filteredRefreshTimer = null;
requestSnapshot(requestedGraphMode);
}, 250);
requestRequiredSnapshot();
}

function prefersFilteredSnapshots() {
return requestedGraphMode === 'filtered' || (requestedGraphMode === 'auto' && graphMeta.largeGraph);
function requestRequiredSnapshot() {
if (snapshotRequestPending) {
return;
}
requestSnapshot(requestedGraphMode);
}

function setGraphMode(mode) {
Expand Down
Loading
Loading