From 7d7c32823fe2359d5d4ded61c6f332e02a4832aa Mon Sep 17 00:00:00 2001
From: Javier Marcos <1271349+javuto@users.noreply.github.com>
Date: Fri, 25 Sep 2026 09:22:36 +0200
Subject: [PATCH] Fix for console and file explorer timeout
---
cmd/api/handlers/console.go | 75 ++++++-------------
cmd/api/handlers/console_test.go | 71 ++++++++----------
cmd/api/handlers/file_explorer.go | 19 ++---
cmd/api/handlers/file_explorer_test.go | 24 +++---
.../features/nodes/NodeConsolePage.test.tsx | 53 ++++++++++++-
.../src/features/nodes/NodeConsolePage.tsx | 4 +
6 files changed, 133 insertions(+), 113 deletions(-)
diff --git a/cmd/api/handlers/console.go b/cmd/api/handlers/console.go
index 3e6f59158..3954ec6a2 100644
--- a/cmd/api/handlers/console.go
+++ b/cmd/api/handlers/console.go
@@ -89,7 +89,7 @@ func (h *HandlersApi) ConsoleSessionCreateHandler(w http.ResponseWriter, r *http
// node's next QueryRead can also switch to fast polling before the
// operator types their first command. A priming failure is non-fatal.
var priming *console.Command
- if primingCmd, primingErr := h.Console.SubmitPrimingCommand(session.ID, h.consolePrimingTimeout(env, node)); primingErr == nil {
+ if primingCmd, primingErr := h.Console.SubmitPrimingCommand(session.ID, h.consolePrimingTimeout(env)); primingErr == nil {
priming = &primingCmd
}
h.auditConsoleVisit(ctx[ctxUser], r, env.ID)
@@ -146,11 +146,7 @@ func (h *HandlersApi) ConsoleCommandCreateHandler(w http.ResponseWriter, r *http
apiErrorResponse(w, "no access", http.StatusForbidden, fmt.Errorf("attempt to use console get by user %s", ctx[ctxUser]))
return
}
- node, ok := h.sessionNode(w, env, session.NodeUUID)
- if !ok {
- return
- }
- command, parsed, err := h.Console.SubmitCommandWithTimeout(session.ID, body.Input, h.consoleCommandTimeout(env, node, preview), body.OsqueryMode)
+ command, parsed, err := h.Console.SubmitCommandWithTimeout(session.ID, body.Input, h.consoleCommandTimeout(env, preview), body.OsqueryMode)
if err != nil {
apiErrorResponse(w, err.Error(), http.StatusBadRequest, err)
return
@@ -257,55 +253,50 @@ func (h *HandlersApi) acceleratedQueryReadSeconds() int64 {
// distributed interval (--distributed_interval = env.QueryInterval), not
// at the accelerated interval. A query expiring before the node's next
// scheduled read is never delivered, which is what forced operators to
-// click refresh repeatedly until acceleration happened to kick in. The
-// last query read recorded for the node makes the next read predictable:
-// last read + configured interval. The wait is whatever remains of that
-// window, capped by maxWarmupWait. Once acceleration is active the last
-// read is fresh, the wait collapses to zero, and the base timeout alone
-// covers delivery on the (now fast) polling cadence.
-func warmupQueryWait(env environments.TLSEnvironment, node nodes.OsqueryNode) time.Duration {
+// click refresh repeatedly until acceleration happened to kick in.
+//
+// The node's recorded last read is deliberately not used to shrink this
+// window. osctrl-tls stamps last_query_read through a batch writer that
+// coalesces check-ins and flushes on --writer-timeout (60s by default),
+// so the stamp can lag the node's real poll by roughly one interval —
+// the same magnitude as the window being predicted. Trusting it made
+// mid-cycle reads look overdue, granted no extra wait, and expired
+// warmup commands undelivered. Reserving the full interval keeps
+// delivery certain for live nodes; once acceleration is active the
+// command completes on the first fast poll and the extra expiration
+// only bounds how long a dead node's command stays pending.
+func warmupQueryWait(env environments.TLSEnvironment) time.Duration {
interval := env.QueryInterval
if interval <= 0 {
interval = environments.DefaultQueryInterval
}
- if node.LastQueryRead.IsZero() {
- // No read recorded yet (new node or rows predating the column):
- // assume the next read can be a full interval away.
- return min(time.Duration(interval)*time.Second, maxWarmupWait)
- }
- wait := time.Duration(interval)*time.Second - time.Since(node.LastQueryRead)
- if wait <= 0 {
- // The scheduled read is overdue — the node is offline or asleep.
- // Keep the base timeout so dead nodes still fail fast.
- return 0
- }
- return min(wait, maxWarmupWait)
+ return min(time.Duration(interval)*time.Second, maxWarmupWait)
}
-func (h *HandlersApi) consoleCommandTimeout(env environments.TLSEnvironment, node nodes.OsqueryNode, parsed console.ParsedCommand) time.Duration {
+func (h *HandlersApi) consoleCommandTimeout(env environments.TLSEnvironment, parsed console.ParsedCommand) time.Duration {
seconds := h.acceleratedQueryReadSeconds()
if parsed.Kind == console.CommandRemote && parsed.Command == "sql" {
timeout := time.Duration(seconds*12) * time.Second
if timeout < time.Minute {
timeout = time.Minute
}
- return timeout + warmupQueryWait(env, node)
+ return timeout + warmupQueryWait(env)
}
- return time.Duration(seconds*2)*time.Second + warmupQueryWait(env, node)
+ return time.Duration(seconds*2)*time.Second + warmupQueryWait(env)
}
// consolePrimingTimeout is the expiration given to the priming metadata
// query. The base is intentionally generous (the accelerated interval
// doubled plus a minute floor) so the priming query stays pending long
-// enough for the next accelerated QueryRead to deliver it, and the warmup
-// wait extends it further while the node is still on its regular polling
-// interval.
-func (h *HandlersApi) consolePrimingTimeout(env environments.TLSEnvironment, node nodes.OsqueryNode) time.Duration {
+// enough for the next accelerated QueryRead to deliver it, and the
+// warmup wait extends it further while the node is still polling at its
+// regular interval.
+func (h *HandlersApi) consolePrimingTimeout(env environments.TLSEnvironment) time.Duration {
timeout := time.Duration(h.acceleratedQueryReadSeconds()*2) * time.Second
if timeout < time.Minute {
timeout = time.Minute
}
- return timeout + warmupQueryWait(env, node)
+ return timeout + warmupQueryWait(env)
}
func osqueryTableSupportsPlatform(table types.OsqueryTable, platform string) bool {
@@ -416,24 +407,6 @@ func (h *HandlersApi) consoleSessionContext(w http.ResponseWriter, r *http.Reque
return env, ctx, session, true
}
-// sessionNode resolves the node an interactive session (console or file
-// explorer) belongs to, so submit paths can size the distributed query
-// expiration from the node's polling state. The session was created
-// against this node; if it no longer resolves, nothing submitted to it
-// can ever be delivered.
-func (h *HandlersApi) sessionNode(w http.ResponseWriter, env environments.TLSEnvironment, nodeUUID string) (nodes.OsqueryNode, bool) {
- node, err := h.Nodes.GetByUUIDEnv(nodeUUID, env.ID)
- if err != nil {
- if errors.Is(err, gorm.ErrRecordNotFound) {
- apiErrorResponse(w, "node not found", http.StatusNotFound, err)
- return nodes.OsqueryNode{}, false
- }
- apiErrorResponse(w, "error getting node", http.StatusInternalServerError, err)
- return nodes.OsqueryNode{}, false
- }
- return node, true
-}
-
func consolePathUint(w http.ResponseWriter, r *http.Request, name string) (uint, bool) {
value := r.PathValue(name)
id, err := strconv.ParseUint(value, 10, strconv.IntSize)
diff --git a/cmd/api/handlers/console_test.go b/cmd/api/handlers/console_test.go
index f9260c89f..84384940b 100644
--- a/cmd/api/handlers/console_test.go
+++ b/cmd/api/handlers/console_test.go
@@ -112,32 +112,26 @@ func TestConsoleSessionCreateDispatchesPrimingCommand(t *testing.T) {
require.True(t, distributed.Hidden)
}
-func TestWarmupQueryWaitSizedToNodePollingState(t *testing.T) {
- env := environments.TLSEnvironment{QueryInterval: 60}
- // Fresh read: nearly the whole interval remains until the next read.
- wait := warmupQueryWait(env, nodes.OsqueryNode{LastQueryRead: time.Now()})
- require.Greater(t, wait, 55*time.Second)
- require.LessOrEqual(t, wait, 60*time.Second)
- // Read 30s ago on a 60s interval: about 30s remain.
- wait = warmupQueryWait(env, nodes.OsqueryNode{LastQueryRead: time.Now().Add(-30 * time.Second)})
- require.Greater(t, wait, 25*time.Second)
- require.LessOrEqual(t, wait, 35*time.Second)
- // Overdue read: no extra wait, so dead nodes still fail fast.
- require.Zero(t, warmupQueryWait(env, nodes.OsqueryNode{LastQueryRead: time.Now().Add(-2 * time.Minute)}))
- // No read recorded yet: assume the next read can be a full interval away.
- require.Equal(t, 60*time.Second, warmupQueryWait(env, nodes.OsqueryNode{}))
+func TestWarmupQueryWaitCoversFullPollInterval(t *testing.T) {
+ // last_query_read is stamped by osctrl-tls's batch writer and can lag
+ // the node's real poll by up to a flush window, so it must not shrink
+ // the warmup window: every warming-up command reserves a full interval
+ // regardless of how fresh or stale the recorded read looks.
+ require.Equal(t, 60*time.Second, warmupQueryWait(environments.TLSEnvironment{QueryInterval: 60}))
+ // A recorded read (fresh or stale) must not change the reservation.
+ require.Equal(t, 60*time.Second, warmupQueryWait(environments.TLSEnvironment{QueryInterval: 60}))
// Unset interval falls back to the environment default.
- require.Equal(t, 60*time.Second, warmupQueryWait(environments.TLSEnvironment{}, nodes.OsqueryNode{}))
+ require.Equal(t, 60*time.Second, warmupQueryWait(environments.TLSEnvironment{}))
// Very long intervals are capped so requests cannot pend unbounded.
- require.Equal(t, maxWarmupWait, warmupQueryWait(environments.TLSEnvironment{QueryInterval: 3600}, nodes.OsqueryNode{LastQueryRead: time.Now()}))
+ require.Equal(t, maxWarmupWait, warmupQueryWait(environments.TLSEnvironment{QueryInterval: 3600}))
}
func TestConsoleSessionCreatePrimingSurvivesNodePollInterval(t *testing.T) {
db, h, env, node := setupConsoleHandlers(t)
require.NoError(t, db.Model(&env).UpdateColumn("query_interval", 60).Error)
- // The node is mid-cycle on its regular distributed interval; the
- // priming query must still be pending when its next read arrives.
- require.NoError(t, db.Model(&node).UpdateColumn("last_query_read", time.Now().Add(-30*time.Second)).Error)
+ // The node's recorded read is mid-cycle on its regular distributed
+ // interval; the priming query must still be pending when its next
+ // read arrives.
before := time.Now()
req := consoleRequest(http.MethodPost, "/console", nil, "alice")
@@ -154,8 +148,9 @@ func TestConsoleSessionCreatePrimingSurvivesNodePollInterval(t *testing.T) {
var distributed queries.DistributedQuery
require.NoError(t, db.Where("name = ?", resp.Priming.DistributedQueryName).First(&distributed).Error)
- require.True(t, distributed.Expiration.After(before.Add(85*time.Second)), "priming must cover the node's next scheduled read")
- require.True(t, distributed.Expiration.Before(before.Add(97*time.Second)))
+ // Base is the minute floor plus the full warmup wait (60s interval).
+ require.True(t, distributed.Expiration.After(before.Add(115*time.Second)), "priming must cover the node's next scheduled read")
+ require.True(t, distributed.Expiration.Before(before.Add(125*time.Second)))
}
func TestConsoleSessionCreateReturnsNodeInfo(t *testing.T) {
@@ -283,9 +278,8 @@ func TestConsoleCommandRejectsSecondInFlightCommand(t *testing.T) {
func TestConsoleCommandExpirationUsesDoubleAcceleratedQueryReadInterval(t *testing.T) {
db, h, env, node := setupConsoleHandlers(t)
require.NoError(t, h.Settings.NewIntegerValue(config.ServiceTLS, settings.AcceleratedSeconds, 7, settings.NoEnvironmentID))
- // The node's scheduled read is overdue, so no warmup wait is added:
- // the expiration is the doubled accelerated interval alone.
- require.NoError(t, db.Model(&node).UpdateColumn("last_query_read", time.Now().Add(-2*time.Minute)).Error)
+ // The expiration is the doubled accelerated interval plus the full
+ // warmup wait (60s default interval).
session, err := h.Console.CreateSession(env, node, "alice")
require.NoError(t, err)
before := time.Now()
@@ -302,19 +296,20 @@ func TestConsoleCommandExpirationUsesDoubleAcceleratedQueryReadInterval(t *testi
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &resp))
var distributed queries.DistributedQuery
require.NoError(t, db.Where("name = ?", resp.Command.DistributedQueryName).First(&distributed).Error)
- require.True(t, distributed.Expiration.After(before.Add(13*time.Second)))
- require.True(t, distributed.Expiration.Before(before.Add(15*time.Second)))
+ require.True(t, distributed.Expiration.After(before.Add(73*time.Second)))
+ require.True(t, distributed.Expiration.Before(before.Add(76*time.Second)))
}
func TestConsoleCommandExpirationSurvivesNodePollIntervalWhileWarmingUp(t *testing.T) {
db, h, env, node := setupConsoleHandlers(t)
require.NoError(t, h.Settings.NewIntegerValue(config.ServiceTLS, settings.AcceleratedSeconds, 7, settings.NoEnvironmentID))
require.NoError(t, db.Model(&env).UpdateColumn("query_interval", 60).Error)
- // The node's last query read was 30s ago on a 60s interval: its next
- // read is up to 30s away. The query must stay pending at least that
- // long, or warmup commands expire undelivered and the operator has to
- // keep clicking refresh until acceleration happens to kick in.
- require.NoError(t, db.Model(&node).UpdateColumn("last_query_read", time.Now().Add(-30*time.Second)).Error)
+ // The node's last recorded read looks overdue, but the batch writer
+ // lag means its next real read can still be a full interval away. The
+ // query must stay pending at least that long, or warmup commands
+ // expire undelivered and the operator has to keep clicking refresh
+ // until acceleration happens to kick in.
+ require.NoError(t, db.Model(&node).UpdateColumn("last_query_read", time.Now().Add(-2*time.Minute)).Error)
session, err := h.Console.CreateSession(env, node, "alice")
require.NoError(t, err)
before := time.Now()
@@ -332,16 +327,16 @@ func TestConsoleCommandExpirationSurvivesNodePollIntervalWhileWarmingUp(t *testi
require.NotNil(t, resp.Command.ExpiresAt, "the response must expose the deadline so clients can wait it out")
var distributed queries.DistributedQuery
require.NoError(t, db.Where("name = ?", resp.Command.DistributedQueryName).First(&distributed).Error)
- require.True(t, distributed.Expiration.After(before.Add(40*time.Second)), "expiration must cover the node's next scheduled read")
- require.True(t, distributed.Expiration.Before(before.Add(47*time.Second)))
+ require.True(t, distributed.Expiration.After(before.Add(73*time.Second)), "expiration must cover the node's next scheduled read despite stale check-in data")
+ require.True(t, distributed.Expiration.Before(before.Add(76*time.Second)))
}
func TestConsoleOsqueryModeSQLUsesLongerExpiration(t *testing.T) {
db, h, env, node := setupConsoleHandlers(t)
require.NoError(t, h.Settings.NewIntegerValue(config.ServiceTLS, settings.AcceleratedSeconds, 5, settings.NoEnvironmentID))
- // Overdue read: no warmup wait, so the expiration is the long sql
- // timeout alone.
- require.NoError(t, db.Model(&node).UpdateColumn("last_query_read", time.Now().Add(-2*time.Minute)).Error)
+ require.NoError(t, db.Model(&env).UpdateColumn("query_interval", 60).Error)
+ // The sql timeout floor is a minute; the full warmup wait (60s) is
+ // added on top so delivery is covered while still warming up.
session, _ := h.Console.CreateSession(env, node, "alice")
before := time.Now()
@@ -357,8 +352,8 @@ func TestConsoleOsqueryModeSQLUsesLongerExpiration(t *testing.T) {
require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &resp))
var distributed queries.DistributedQuery
require.NoError(t, db.Where("name = ?", resp.Command.DistributedQueryName).First(&distributed).Error)
- require.True(t, distributed.Expiration.After(before.Add(59*time.Second)))
- require.True(t, distributed.Expiration.Before(before.Add(61*time.Second)))
+ require.True(t, distributed.Expiration.After(before.Add(119*time.Second)))
+ require.True(t, distributed.Expiration.Before(before.Add(122*time.Second)))
}
func TestConsoleCommandRejectsNonAdminSessionOwner(t *testing.T) {
diff --git a/cmd/api/handlers/file_explorer.go b/cmd/api/handlers/file_explorer.go
index 632eac307..34f968959 100644
--- a/cmd/api/handlers/file_explorer.go
+++ b/cmd/api/handlers/file_explorer.go
@@ -10,7 +10,6 @@ import (
"github.com/jmpsec/osctrl/pkg/environments"
"github.com/jmpsec/osctrl/pkg/fileexplorer"
- "github.com/jmpsec/osctrl/pkg/nodes"
"github.com/jmpsec/osctrl/pkg/types"
"github.com/jmpsec/osctrl/pkg/utils"
"gorm.io/gorm"
@@ -61,7 +60,7 @@ func (h *HandlersApi) FileExplorerSessionCreateHandler(w http.ResponseWriter, r
// before the operator expands the first directory. Non-fatal on
// failure.
var priming *fileexplorer.Request
- if primingReq, primingErr := h.FileExplorer.SubmitPrimingRequest(session.ID, h.fileExplorerRequestTimeout(env, node)); primingErr == nil {
+ if primingReq, primingErr := h.FileExplorer.SubmitPrimingRequest(session.ID, h.fileExplorerRequestTimeout(env)); primingErr == nil {
priming = &primingReq
}
h.auditFileExplorerAction(ctx[ctxUser], "file explorer session", r, env.ID)
@@ -107,11 +106,7 @@ func (h *HandlersApi) FileExplorerListHandler(w http.ResponseWriter, r *http.Req
if !ok {
return
}
- node, ok := h.sessionNode(w, env, session.NodeUUID)
- if !ok {
- return
- }
- request, err := h.FileExplorer.ListDirectory(session.ID, path, h.fileExplorerRequestTimeout(env, node))
+ request, err := h.FileExplorer.ListDirectory(session.ID, path, h.fileExplorerRequestTimeout(env))
if err != nil {
apiErrorResponse(w, err.Error(), http.StatusBadRequest, err)
return
@@ -129,11 +124,7 @@ func (h *HandlersApi) FileExplorerStatHandler(w http.ResponseWriter, r *http.Req
if !ok {
return
}
- node, ok := h.sessionNode(w, env, session.NodeUUID)
- if !ok {
- return
- }
- request, err := h.FileExplorer.StatPath(session.ID, path, h.fileExplorerRequestTimeout(env, node))
+ request, err := h.FileExplorer.StatPath(session.ID, path, h.fileExplorerRequestTimeout(env))
if err != nil {
apiErrorResponse(w, err.Error(), http.StatusBadRequest, err)
return
@@ -260,8 +251,8 @@ func fileExplorerRequestPath(w http.ResponseWriter, r *http.Request) (string, bo
// once the node polls at the accelerated interval; the warmup wait keeps
// the query alive until the node's next regularly scheduled read while
// acceleration has not kicked in yet.
-func (h *HandlersApi) fileExplorerRequestTimeout(env environments.TLSEnvironment, node nodes.OsqueryNode) time.Duration {
- return time.Duration(h.acceleratedQueryReadSeconds()*2)*time.Second + warmupQueryWait(env, node)
+func (h *HandlersApi) fileExplorerRequestTimeout(env environments.TLSEnvironment) time.Duration {
+ return time.Duration(h.acceleratedQueryReadSeconds()*2)*time.Second + warmupQueryWait(env)
}
func (h *HandlersApi) auditFileExplorerAction(user, action string, r *http.Request, envID uint) {
diff --git a/cmd/api/handlers/file_explorer_test.go b/cmd/api/handlers/file_explorer_test.go
index 6d61243de..0c40c885c 100644
--- a/cmd/api/handlers/file_explorer_test.go
+++ b/cmd/api/handlers/file_explorer_test.go
@@ -110,10 +110,9 @@ func TestFileExplorerSessionCreateDispatchesPrimingRequest(t *testing.T) {
func TestFileExplorerSessionCreatePrimingSurvivesNodePollInterval(t *testing.T) {
db, h, env, node := setupFileExplorerHandlers(t)
require.NoError(t, db.Model(&env).UpdateColumn("query_interval", 60).Error)
- // The node is mid-cycle on its regular distributed interval; the
- // priming query (and the first directory listing that follows) must
- // still be pending when its next read arrives.
- require.NoError(t, db.Model(&node).UpdateColumn("last_query_read", time.Now().Add(-30*time.Second)).Error)
+ // The node's recorded read is mid-cycle on its regular distributed
+ // interval; the priming query (and the first directory listing that
+ // follows) must still be pending when its next read arrives.
before := time.Now()
req := fileExplorerRequest(http.MethodPost, "/file-explorer", nil, "alice")
@@ -130,14 +129,19 @@ func TestFileExplorerSessionCreatePrimingSurvivesNodePollInterval(t *testing.T)
var distributed queries.DistributedQuery
require.NoError(t, db.Where("name = ?", resp.Priming.DistributedQueryName).First(&distributed).Error)
- require.True(t, distributed.Expiration.After(before.Add(35*time.Second)), "priming must cover the node's next scheduled read")
- require.True(t, distributed.Expiration.Before(before.Add(47*time.Second)))
+ // Base is the doubled accelerated interval (5s default) plus the full
+ // warmup wait (60s interval).
+ require.True(t, distributed.Expiration.After(before.Add(65*time.Second)), "priming must cover the node's next scheduled read")
+ require.True(t, distributed.Expiration.Before(before.Add(75*time.Second)))
}
func TestFileExplorerListRequestSurvivesNodePollInterval(t *testing.T) {
db, h, env, node := setupFileExplorerHandlers(t)
require.NoError(t, db.Model(&env).UpdateColumn("query_interval", 60).Error)
- require.NoError(t, db.Model(&node).UpdateColumn("last_query_read", time.Now().Add(-30*time.Second)).Error)
+ // A recorded read that looks overdue must not shrink the warmup
+ // window: the batch writer stamps last_query_read late, so the next
+ // real read can still be a full interval away.
+ require.NoError(t, db.Model(&node).UpdateColumn("last_query_read", time.Now().Add(-2*time.Minute)).Error)
session, err := h.FileExplorer.CreateSession(env, node, "alice")
require.NoError(t, err)
before := time.Now()
@@ -155,8 +159,10 @@ func TestFileExplorerListRequestSurvivesNodePollInterval(t *testing.T) {
var distributed queries.DistributedQuery
require.NoError(t, db.Where("name = ?", request.DistributedQueryName).First(&distributed).Error)
- require.True(t, distributed.Expiration.After(before.Add(35*time.Second)), "the listing must cover the node's next scheduled read")
- require.True(t, distributed.Expiration.Before(before.Add(47*time.Second)))
+ // Base is the doubled accelerated interval (5s default) plus the full
+ // warmup wait (60s interval), despite the stale recorded read.
+ require.True(t, distributed.Expiration.After(before.Add(65*time.Second)), "the listing must cover the node's next scheduled read despite stale check-in data")
+ require.True(t, distributed.Expiration.Before(before.Add(75*time.Second)))
}
func TestFileExplorerPrimingResultsReturnsOsqueryInfoRows(t *testing.T) {
diff --git a/frontend/src/features/nodes/NodeConsolePage.test.tsx b/frontend/src/features/nodes/NodeConsolePage.test.tsx
index a4c442b2c..9d47f33fb 100644
--- a/frontend/src/features/nodes/NodeConsolePage.test.tsx
+++ b/frontend/src/features/nodes/NodeConsolePage.test.tsx
@@ -1,5 +1,5 @@
import { QueryClient, QueryClientProvider } from '@tanstack/react-query';
-import { fireEvent, render, screen, waitFor } from '@testing-library/react';
+import { act, fireEvent, render, screen, waitFor } from '@testing-library/react';
import type { ReactNode } from 'react';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { NodeConsolePanel } from './NodeConsolePage';
@@ -13,6 +13,14 @@ const consoleApi = vi.hoisted(() => ({
getConsoleCommandResults: vi.fn(),
}));
+const liveUpdates = vi.hoisted(() => ({
+ readEvents: vi.fn(),
+ getFeatures: vi.fn(),
+}));
+
+vi.mock('$/api/events', () => ({ readEvents: liveUpdates.readEvents }));
+vi.mock('$/api/features', () => ({ getFeatures: liveUpdates.getFeatures }));
+
vi.mock('@tanstack/react-router', () => ({
Link: ({ children }: { children: ReactNode }) => {children},
useNavigate: () => vi.fn(),
@@ -30,6 +38,8 @@ vi.mock('$/api/console', () => ({
describe('NodeConsolePanel', () => {
beforeEach(() => {
vi.clearAllMocks();
+ liveUpdates.getFeatures.mockResolvedValue({ events: true, event_topics: ['console'] });
+ liveUpdates.readEvents.mockImplementation(() => new Promise(() => undefined));
consoleApi.createConsoleSession.mockResolvedValue({
session: makeSession(),
history: [],
@@ -86,6 +96,47 @@ describe('NodeConsolePanel', () => {
expect(input).toHaveFocus();
});
+ it('refetches the priming command when a live update marks it changed', async () => {
+ consoleApi.createConsoleSession.mockResolvedValue({
+ session: makeSession(),
+ history: [],
+ priming: makeCommand({ id: 99, input: 'select version from osquery_info', priming: true }),
+ });
+ let receive: ((event: { event: string; data: Record }) => void) | undefined;
+ liveUpdates.readEvents.mockImplementation((_env: string, _topics: string[], _signal: AbortSignal, onEvent: typeof receive) => {
+ receive = onEvent;
+ return new Promise(() => undefined);
+ });
+ let primingReads = 0;
+ consoleApi.getConsoleCommand.mockImplementation((_env: string, _sessionID: number, commandID: number) => {
+ if (commandID === 99) {
+ primingReads += 1;
+ return Promise.resolve(makeCommand({ id: 99, input: 'select version from osquery_info', priming: true, status: 'queued' }));
+ }
+ return Promise.resolve(makeCommand({ status: 'queued', input: 'ps' }));
+ });
+
+ renderConsole();
+ const input = await screen.findByLabelText(/console input/i);
+ await waitFor(() => expect(input).not.toBeDisabled());
+ // The session-scoped console stream must be connected before events
+ // can drive the priming poll.
+ await waitFor(() => expect(liveUpdates.readEvents).toHaveBeenCalled());
+ const sessionScoped = liveUpdates.readEvents.mock.calls.some(([, topics, _signal, , options]) =>
+ Array.isArray(topics) && topics.includes('console') && options?.consoleSessionId === 1,
+ );
+ expect(sessionScoped).toBe(true);
+ const readsBefore = primingReads;
+ act(() => {
+ receive?.({ event: 'stream.ready', data: { environment_uuid: 'env-uuid' } });
+ receive?.({
+ event: 'resource.changed',
+ data: { schema_version: 1, environment_uuid: 'env-uuid', topic: 'console', session_id: 1, resource_id: 99, change: 'results' },
+ });
+ });
+ await waitFor(() => expect(primingReads).toBeGreaterThan(readsBefore));
+ });
+
it('keeps the console session fresh while open', async () => {
const intervals: Array<() => void> = [];
const originalSetInterval = window.setInterval.bind(window);
diff --git a/frontend/src/features/nodes/NodeConsolePage.tsx b/frontend/src/features/nodes/NodeConsolePage.tsx
index ec2ea6405..97c887588 100644
--- a/frontend/src/features/nodes/NodeConsolePage.tsx
+++ b/frontend/src/features/nodes/NodeConsolePage.tsx
@@ -70,6 +70,10 @@ export function NodeConsolePanel({ env, uuid }: { env: string; uuid: string }) {
useCallback(({ resourceId }) => {
if (!sessionID) return;
void qc.invalidateQueries({ queryKey: ['console-command', env, sessionID, resourceId], exact: true });
+ // The priming poll uses its own query key; without invalidation the
+ // "warming" badge and header metadata lag up to a reconcile interval
+ // after the node delivers the priming results.
+ void qc.invalidateQueries({ queryKey: ['console-priming', env, sessionID], exact: false });
}, [env, qc, sessionID]),
);