Integration Name
Tenable Vulnerability Management [packages/tenable_io]
Dataset Name
All
Integration Version
4.14.0
Agent Version
9.5.3
Agent Output Type
logstash
Elasticsearch Version
9.5.2
OS Version and Architecture
Red Hat Enterprise Linux 9.6
Software/API Version
N/A
Error Message
No error message was displayed, not with debug logging enabled, not with CEL tracing enabled.
Event Original
N/A
There are no problems with the events being generated by the integration.
What did you do?
Applied stock integration code with three streams enabled (asset, vulnerability, plugin). Below you will find the code we are using right now, with a small diagnostic patch added to test the theory that API throttling was causing our problem, we we did see throttling messages from the API server. However, the patch did not fix the issue.
inputs:
- id: cel-tenable_io-fc7810d5-2605-4f13-a702-fcbddce1630f
type: cel
data_stream.namespace: default
streams:
- config_version: 2
interval: ${env.TENABLE_IO_ASSET_INTERVAL}
resource.tracer:
enabled: false
filename: ../../logs/cel/http-request-trace-*.ndjson
maxbackups: 5
resource.ssl: null
resource.timeout: ${env.TENABLE_IO_ASSET_RESOURCE_TIMEOUT}
resource.url: ${env.CEL_RESOURCE_URL}
state:
access_key: ${env.ACCESS_KEY}
secret_key: ${env.SECRET_KEY}
batch_size: 100
initial_interval: ${env.TENABLE_IO_ASSET_INITIAL_INTERVAL}
export_status_timeout: ${env.TENABLE_IO_ASSET_EXPORT_STATUS_TIMEOUT}
export_status_poll_interval: ${env.TENABLE_IO_ASSET_RESOURCE_TIMEOUT}
redact:
fields:
- access_key
- secret_key
max_executions: 2000
program: |-
// The asset export is a small state machine that makes one HTTP
// request per evaluation. state.export holds the active export job
// ({uuid, created, expires}); state.chunks and state.next track the chunk
// downloads. Everything except the cursor is in-memory only and is lost on
// restart, so the cursor (last_event_ts) is advanced only once every
// chunk of an export has been downloaded. A cancelled, failed, expired or
// orphaned export is therefore retried from the same point in time at the
// next interval instead of silently skipping its window.
//
// Phases, selected in this order: "create" when there is no active export,
// which requests a new export since the cursor; "download" when the active
// export has chunks left to fetch, so an interrupted download resumes while
// Tenable still has the chunks; "cancel" when the active export is still
// being built past export_status_timeout, which cancels and drops it so the
// next interval starts afresh; otherwise "status", which polls the export
// job. QUEUED and PROCESSING wait export_status_poll_interval before the
// next poll, while CANCELLED and ERROR are
// terminal and reported as errors.
{
"X-ApiKeys": ["accessKey=" + state.access_key + ";secretKey=" + state.secret_key],
"User-Agent": ["Integration/1.0 (Elastic; Tenable.io; Build/3.0.0)"],
}.as(headers,
(
!has(state.?export.uuid) ?
"create"
: (size(state.?chunks.orValue([])) > 0) ?
"download"
: (int(now) > int(state.export.expires)) ?
"cancel"
:
"status"
).as(phase,
state.with(
(phase == "create") ?
// Create: request an export for everything updated since the cursor.
post_request(
state.url.trim_right("/") + "/assets/export",
"application/json",
{
"chunk_size": state.batch_size,
"filters": {
"updated_at": int(state.?cursor.last_event_ts.orValue(int(now - duration(state.initial_interval)))),
},
}.encode_json()
).with({"Header": headers}).do_request().as(resp,
(resp.StatusCode == 200) ?
{
"export": {
"uuid": resp.Body.decode_json().export_uuid,
"created": int(now),
"expires": int(now + duration(state.export_status_timeout)),
},
"chunks": [],
"next": 0,
"events": [{"retry": true}],
"want_more": true,
}
:
{
"events": {
"error": {
"code": string(resp.StatusCode),
"id": string(resp.Status),
"message": "POST /assets/export: " + ((size(resp.Body) != 0) ? string(resp.Body) : string(resp.Status)),
},
},
"want_more": false,
}
)
: (phase == "cancel") ?
// Cancel: the export did not finish in time. Drop it whatever the
// API says; the cursor is untouched so the next interval retries.
post_request(
state.url.trim_right("/") + "/assets/export/" + state.export.uuid + "/cancel",
"application/json",
""
).with({"Header": headers}).do_request().as(resp,
{
"export": {},
"chunks": [],
"next": 0,
"events": {
"error": {
"code": "EXPIRED",
"id": state.export.uuid,
"message": "POST /assets/export/" + state.export.uuid + "/cancel: export did not finish within " + state.export_status_timeout + (
(resp.StatusCode == 200) ?
""
:
"; cancel failed: " + ((size(resp.Body) != 0) ? string(resp.Body) : string(resp.Status))
),
},
},
"want_more": false,
}
)
: (phase == "download") ?
// Download: fetch the next chunk. Once the last chunk is in, the
// export is complete and the cursor moves to the export's creation
// time.
request(
"GET",
state.url.trim_right("/") + "/assets/export/" + state.export.uuid + "/chunks/" + string(int(state.chunks[int(state.next)]))
).with({"Header": headers}).do_request().as(resp,
(resp.StatusCode == 200) ?
resp.Body.decode_json().as(body,
(int(state.next) + 1 < size(state.chunks)).as(more,
{
"events": (body != null && size(body) > 0) ?
dyn(body.map(e, {"message": e.encode_json()}))
:
dyn([{"retry": true}]),
"export": more ? state.export : {},
"chunks": more ? state.chunks : [],
"next": more ? (int(state.next) + 1) : 0,
?"cursor": more ? optional.none() : optional.of({"last_event_ts": int(state.export.created)}),
"want_more": more,
}
)
)
:
// The chunk is gone (Tenable expires chunks after 24h) or keeps
// failing after the client's own retries: drop the export. The
// cursor is untouched, so the next interval creates a fresh export
// from the same point.
{
"export": {},
"chunks": [],
"next": 0,
"events": {
"error": {
"code": string(resp.StatusCode),
"id": string(resp.Status),
"message": "GET /assets/export/" + state.export.uuid + "/chunks/" + string(int(state.chunks[int(state.next)])) + ": " + ((size(resp.Body) != 0) ? string(resp.Body) : string(resp.Status)),
},
},
"want_more": false,
}
)
:
// Status: poll the export job.
request(
"GET",
state.url.trim_right("/") + "/assets/export/" + state.export.uuid + "/status"
).with({"Header": headers}).do_request().as(resp,
(resp.StatusCode == 200) ?
resp.Body.decode_json().as(body,
(body.?status.orValue("") == "FINISHED" && size(body.?chunks_available.orValue([])) > 0) ?
{
"chunks": body.chunks_available,
"next": 0,
"events": [{"retry": true}],
"want_more": true,
}
: (body.?status.orValue("") == "FINISHED") ?
// Finished with nothing to download: the export is complete.
{
"export": {},
"chunks": [],
"next": 0,
"cursor": {"last_event_ts": int(state.export.created)},
"events": [{"retry": true}],
"want_more": false,
}
: (body.?status.orValue("") in ["CANCELLED", "ERROR"]) ?
// Terminal failure: report it and drop the export. The cursor
// is untouched so the next interval retries the same window.
{
"export": {},
"chunks": [],
"next": 0,
"events": {
"error": {
"code": body.status,
"id": state.export.uuid,
"message": "GET /assets/export/" + state.export.uuid + "/status: export " + body.status + ((body.?reason.orValue(null) != null) ? (": " + string(body.reason)) : ""),
},
},
"want_more": false,
}
:
// QUEUED or PROCESSING: ask the input to wait before the next
// evaluation instead of polling back-to-back.
{
"events": [{"retry": true}],
"want_more": true,
"rate_limit": {
"rate": 0,
"burst": 1,
"next": "inf",
"reset": string(now + duration(state.export_status_poll_interval)),
},
}
)
:
// Keep the export so the next interval polls it again (or
// cancels it once it has expired).
{
"events": {
"error": {
"code": string(resp.StatusCode),
"id": string(resp.Status),
"message": "GET /assets/export/" + state.export.uuid + "/status: " + ((size(resp.Body) != 0) ? string(resp.Body) : string(resp.Status)),
},
},
"want_more": false,
}
)
)
)
)
tags:
- preserve_original_event
- preserve_duplicate_custom_fields
- forwarded
- tenable_io-asset
publisher_pipeline.disable_host: true
processors:
- drop_event.when.equals.retry: true
data_stream:
type: logs
dataset: tenable_io.asset
- config_version: 2
interval: ${env.TENABLE_IO_PLUGIN_INTERVAL}
resource.tracer:
enabled: false
filename: ../../logs/cel/http-request-trace-*.ndjson
maxbackups: 5
resource.ssl: null
resource.timeout: ${env.TENABLE_IO_PLUGIN_RESOURCE_TIMEOUT}
resource.url: ${env.CEL_RESOURCE_URL}
state:
access_key: ${env.ACCESS_KEY}
secret_key: ${env.SECRET_KEY}
batch_size: 1000
next_page: 1
want_more: false
redact:
fields:
- access_key
- secret_key
max_executions: 2000
program: |
request("GET",
state.url.trim_right("/") + "/plugins/plugin?last_updated=" + (
!state.?want_more.orValue(false) ?
state.?cursor.last_update_time.orValue(string(timestamp(0)))
:
state.?cursor.first_update_time.orValue(null)
) + "&page=" + string(state.next_page) + "&size=" + string(state.batch_size)
).with({
"Header":{
"X-ApiKeys": ["accessKey=" + state.access_key + ";secretKey=" + state.secret_key],
"User-Agent": ["Integration/1.0 (Elastic; Tenable.io; Build/3.0.0)"]
},
}).do_request().as(resp,
resp.StatusCode == 200 ?
bytes(resp.Body).decode_json().as(body,
(int(state.?count.orValue(0))+int(body.size)).as(count,
(count < int(body.total_count)).as(want_more, {
"events": has(body.?data.plugin_details) ?
body.data.plugin_details.map(e, { "message": e.encode_json() })
:
[],
"cursor": {
"last_update_time": (
has(body.?data.plugin_details) && body.data.plugin_details.size() > 0 ?
([?state.?cursor.last_update_time] + body.data.plugin_details.map(e,
e.attributes.plugin_modification_date
)).map(t,
timestamp(t)
).max()
:
state.?cursor.last_update_time.orValue(string(timestamp(0)))
),
"first_update_time": (
has(state.?cursor.first_update_time) && has(body.data.plugin_details) ?
(
want_more ?
state.cursor.first_update_time
:
state.cursor.last_update_time
)
:
string(timestamp(0))
),
},
"next_page": want_more ? int(state.next_page)+1 : 1,
"want_more": want_more,
"count": count,
"access_key": state.access_key,
"secret_key": state.secret_key,
"batch_size": state.batch_size,
})
)
)
:
{
"events": {
"error": {
"code": string(resp.StatusCode),
"id": string(resp.Status),
"message": "GET:"+(
size(resp.Body) != 0 ?
string(resp.Body)
:
string(resp.Status) + ' (' + string(resp.StatusCode) + ')'
),
},
},
"want_more": false,
"next_page": 1,
"batch_size": state.batch_size,
"access_key": state.access_key,
"secret_key": state.secret_key,
}
)
tags:
- preserve_original_event
- preserve_duplicate_custom_fields
- forwarded
- tenable_io-plugin
publisher_pipeline.disable_host: true
data_stream:
type: logs
dataset: tenable_io.plugin
- config_version: 2
interval: ${env.TENABLE_IO_VULNERABILITY_INTERVAL}
resource.tracer:
enabled: false
filename: ../../logs/cel/http-request-trace-*.ndjson
maxbackups: 5
resource.ssl: null
resource.timeout: ${env.TENABLE_IO_VULNERABILITY_RESOURCE_TIMEOUT}
resource.url: ${env.CEL_RESOURCE_URL}
state:
access_key: ${env.ACCESS_KEY}
secret_key: ${env.SECRET_KEY}
batch_size: 50
initial_interval: ${env.TENABLE_IO_VULNERABILITY_INITIAL_INTERVAL}
export_status_timeout: ${env.TENABLE_IO_VULNERABILITY_EXPORT_STATUS_TIMEOUT}
export_status_poll_interval: ${env.TENABLE_IO_VULNERABILITY_RESOURCE_TIMEOUT}
severity_level:
- critical
- high
- medium
- low
- info
redact:
fields:
- access_key
- secret_key
max_executions: 2000
program: |-
// The vulnerability export is a small state machine that makes one HTTP
// request per evaluation. state.export holds the active export job
// ({uuid, created, expires}); state.chunks and state.next track the chunk
// downloads. Everything except the cursor is in-memory only and is lost on
// restart, so the cursor (last_event_time) is advanced only once every
// chunk of an export has been downloaded. A cancelled, failed, expired or
// orphaned export is therefore retried from the same point in time at the
// next interval instead of silently skipping its window.
//
// Phases, selected in this order: "create" when there is no active export,
// which requests a new export since the cursor; "download" when the active
// export has chunks left to fetch, so an interrupted download resumes while
// Tenable still has the chunks; "cancel" when the active export is still
// being built past export_status_timeout, which cancels and drops it so the
// next interval starts afresh; otherwise "status", which polls the export
// job. QUEUED and PROCESSING wait export_status_poll_interval before the
// next poll, while CANCELLED, ERROR and FINISHED with failed chunks are
// terminal and reported as errors.
{
"X-ApiKeys": ["accessKey=" + state.access_key + ";secretKey=" + state.secret_key],
"User-Agent": ["Integration/1.0 (Elastic; Tenable.io; Build/3.0.0)"],
}.as(headers,
(
!has(state.?export.uuid) ?
"create"
: (size(state.?chunks.orValue([])) > 0) ?
"download"
: (int(now) > int(state.export.expires)) ?
"cancel"
:
"status"
).as(phase,
state.with(
(phase == "create") ?
// Create: request an export for everything updated since the cursor.
post_request(
state.url.trim_right("/") + "/vulns/export",
"application/json",
{
"num_assets": state.batch_size,
"filters": {
"since": int(state.?cursor.last_event_time.orValue(int(now - duration(state.initial_interval)))),
"state": ["open", "reopened", "fixed"],
"severity": state.?severity_level.orValue([]),
},
}.encode_json()
).with({"Header": headers}).do_request().as(resp,
(resp.StatusCode == 200) ?
{
"export": {
"uuid": resp.Body.decode_json().export_uuid,
"created": int(now),
"expires": int(now + duration(state.export_status_timeout)),
},
"chunks": [],
"next": 0,
"events": [{"retry": true}],
"want_more": true,
}
:
{
"events": {
"error": {
"code": string(resp.StatusCode),
"id": string(resp.Status),
"message": "POST /vulns/export: " + ((size(resp.Body) != 0) ? string(resp.Body) : string(resp.Status)),
},
},
"want_more": false,
}
)
: (phase == "cancel") ?
// Cancel: the export did not finish in time. Drop it whatever the
// API says; the cursor is untouched so the next interval retries.
post_request(
state.url.trim_right("/") + "/vulns/export/" + state.export.uuid + "/cancel",
"application/json",
""
).with({"Header": headers}).do_request().as(resp,
{
"export": {},
"chunks": [],
"next": 0,
"events": {
"error": {
"code": "EXPIRED",
"id": state.export.uuid,
"message": "POST /vulns/export/" + state.export.uuid + "/cancel: export did not finish within " + state.export_status_timeout + (
(resp.StatusCode == 200) ?
""
:
"; cancel failed: " + ((size(resp.Body) != 0) ? string(resp.Body) : string(resp.Status))
),
},
},
"want_more": false,
}
)
: (phase == "download") ?
// Download: fetch the next chunk. Once the last chunk is in, the
// export is complete and the cursor moves to the export's creation
// time.
request(
"GET",
state.url.trim_right("/") + "/vulns/export/" + state.export.uuid + "/chunks/" + string(int(state.chunks[int(state.next)]))
).with({"Header": headers}).do_request().as(resp,
(resp.StatusCode == 200) ?
resp.Body.decode_json().as(body,
(int(state.next) + 1 < size(state.chunks)).as(more,
{
"events": (body != null && size(body) > 0) ?
dyn(body.map(e, {"message": e.encode_json()}))
:
dyn([{"retry": true}]),
"export": more ? state.export : {},
"chunks": more ? state.chunks : [],
"next": more ? (int(state.next) + 1) : 0,
?"cursor": more ? optional.none() : optional.of({"last_event_time": int(state.export.created)}),
"want_more": more,
}
)
)
:
// The chunk is gone (Tenable expires chunks after 24h) or keeps
// failing after the client's own retries: drop the export. The
// cursor is untouched, so the next interval creates a fresh export
// from the same point.
{
"export": {},
"chunks": [],
"next": 0,
"events": {
"error": {
"code": string(resp.StatusCode),
"id": string(resp.Status),
"message": "GET /vulns/export/" + state.export.uuid + "/chunks/" + string(int(state.chunks[int(state.next)])) + ": " + ((size(resp.Body) != 0) ? string(resp.Body) : string(resp.Status)),
},
},
"want_more": false,
}
)
:
// Status: poll the export job.
request(
"GET",
state.url.trim_right("/") + "/vulns/export/" + state.export.uuid + "/status"
).with({"Header": headers}).do_request().as(resp,
(resp.StatusCode == 200) ?
resp.Body.decode_json().as(body,
(body.?status.orValue("") == "FINISHED" && size(body.?chunks_failed.orValue([])) > 0) ?
// Some chunks failed. Tenable says to submit the export again,
// so drop it and leave the cursor alone; the next interval
// retries the same window.
{
"export": {},
"chunks": [],
"next": 0,
"events": {
"error": {
"code": "CHUNKS_FAILED",
"id": state.export.uuid,
"message": "GET /vulns/export/" + state.export.uuid + "/status: " + string(size(body.chunks_failed)) + " chunk(s) failed" + ((body.?reason.orValue(null) != null) ? (": " + string(body.reason)) : ""),
},
},
"want_more": false,
}
: (body.?status.orValue("") == "FINISHED" && size(body.?chunks_available.orValue([])) > 0) ?
{
"chunks": body.chunks_available,
"next": 0,
"events": [{"retry": true}],
"want_more": true,
}
: (body.?status.orValue("") == "FINISHED") ?
// Finished with nothing to download: the export is complete.
{
"export": {},
"chunks": [],
"next": 0,
"cursor": {"last_event_time": int(state.export.created)},
"events": [{"retry": true}],
"want_more": false,
}
: (body.?status.orValue("") in ["CANCELLED", "ERROR"]) ?
// Terminal failure: report it and drop the export. The cursor
// is untouched so the next interval retries the same window.
{
"export": {},
"chunks": [],
"next": 0,
"events": {
"error": {
"code": body.status,
"id": state.export.uuid,
"message": "GET /vulns/export/" + state.export.uuid + "/status: export " + body.status + ((body.?reason.orValue(null) != null) ? (": " + string(body.reason)) : ""),
},
},
"want_more": false,
}
:
// QUEUED or PROCESSING: ask the input to wait before the next
// evaluation instead of polling back-to-back.
//
// DIAGNOSTIC PATCH, local to this spec's cached policy snapshot only: the
// original "rate": 0 / "reset": <timestamp> form exercises the input's
// zero-rate special case (input.go's handleRateLimit + waitForRateLimit,
// an RFC3339 string round-trip and a manually-scheduled rate.Limiter
// transition). That is the leading suspect for an observed hang: the
// input stops issuing any further requests -- not just this one, all of
// them -- after some minutes of otherwise-successful polling, with no
// error and a stable goroutine count (consistent with a goroutine parked
// forever on an unbounded wait, not a crash). This replaces it with a
// plain positive rate limit (1 request per export_status_poll_interval),
// which takes the far more common, better-exercised code path
// (rate.Limiter.SetLimit/SetBurst) and avoids the reset-timestamp parse
// and the zero-rate scheduling entirely. If the stream keeps advancing
// past where it previously froze, that confirms the zero-rate path as the
// defect without needing a goroutine dump. Revert once confirmed either
// way -- this is not a production fix, just a controlled A/B test.
{
"events": [{"retry": true}],
"want_more": true,
"rate_limit": {
"rate": 1.0 / double(duration(state.export_status_poll_interval).getSeconds()),
"burst": 1,
},
}
)
:
// Keep the export so the next interval polls it again (or
// cancels it once it has expired).
{
"events": {
"error": {
"code": string(resp.StatusCode),
"id": string(resp.Status),
"message": "GET /vulns/export/" + state.export.uuid + "/status: " + ((size(resp.Body) != 0) ? string(resp.Body) : string(resp.Status)),
},
},
"want_more": false,
}
)
)
)
)
tags:
- preserve_original_event
- preserve_duplicate_custom_fields
- forwarded
- tenable_io-vulnerability
publisher_pipeline.disable_host: true
processors:
- drop_event.when.equals.retry: true
data_stream:
type: logs
dataset: tenable_io.vulnerability
processors:
- add_fields:
fields:
pipelines: agent-tenable-io
target: yale
- script:
lang: javascript
source: |
function process(event) {
var now = new Date().toISOString();
event.Put("event.ingested", now);
return event;
}
What did you see?
Only one data stream collects (asset). Remaining data streams never collect data until the other data streams are disabled. i.e to collect "vulnerability" (third in the array or streams we enabled), we need to disable both "asset" and "plugin" streams.
What did you expect to see?
Three data streams collecting consistently (asset, plugin, vulnerability)
Anything else?
Yes... we applied Claude Code to diagnose the problem, and on review the agent suggested that the problem is that the API server URL is shared across all three streams, which triggers a bug in the CEL input cursor framework. The evidence presented seems strong, so I will provide the agent diagnostics below. Before I do that I will note that there are two acceptable fixes from my standpoint:
- Give each data stream a distinct
resource.url by moving the dataset-specific path into the templated URL rather than building it in the CEL program from a shared base. This is the easiest fix as it is specific to this integration, but I think we likely would see the same issue on other multi-stream integrations that share the same URL. I cannot verify this as all the CEL-based integrations that my team has used so far all use a different URL per-stream (see: Abnormal Security, Crowdstrike, Wiz).
- Fix the upstream deadlock in the CEL input. This approach may prevent future bugs with other integrations as well, but would require more testing across all existing CEL-based integrations.
The following the Claude's own root cause analysis and proposed fix. I lack the golang background to evaluate the accuracy of the report fully, but I can tell you that the AI identified the potential bug in the source first, and then verified though a stack analysis gathered from a running agent using elastic-agent diagnostics collect, so it's claims feel defensible to me.
Root cause
The package's CEL program is not involved. tenable_io's asset program signals completion of an
export cycle with "want_more": false once its cursor advances, which is correct. The exclusivity
is two layers below, in the input-cursor framework.
1. The lock key is the URL alone (x-pack/filebeat/input/cel/input_manager.go):
type source struct{ cfg config }
func (s *source) Name() string { return s.cfg.Resource.URL.String() }
Source.Name() is what the framework uses to decide whether two configured streams are the same
resource. For asset and vulnerability it returns the same string,
https://cloud.tenable.com, for the reason in point 4.
config.DataStream carries the stream's own data_stream.dataset on the same struct, populated in
cursorConfigure (src.cfg.DataStream = dataStreamName(cfg)) and already used elsewhere in the
input's lifecycle for status reporting. Name() does not read it.
2. The lock is held for the input's lifetime (filebeat/input/v2/input-cursor/):
// manager.go
// The InputManager ensures that only one input can be active for a unique source.
// If two inputs have overlapping sources, both can still collect data, but
// only one input will collect from the common source.
// input.go, managedInput.runSource
resourceKey := inp.createSourceID(source) // "<type>::<Source.Name()>"
resource, err := lock(ctx, store, resourceKey) // blocks until free
defer releaseResource(resource)
return inp.input.Run(ctx, source, cursor, publisher)
3. Run() does not return during normal operation
(x-pack/filebeat/input/cel/input.go):
func (i input) run(env v2.Context, src *source, ...) error {
return periodically(ctx, cfg.Interval, s.runCycle)
}
func periodically(ctx context.Context, each time.Duration, fn func(context.Context) error) error {
err := fn(ctx)
if err != nil { return err }
return timed.Periodic(ctx, each, func() error { return fn(ctx) }) // blocks until ctx.Done()
}
timed.Periodic blocks until the input's context ends, re-running runCycle every interval. A
CEL program finishing its export window ends one runCycle call, not the surrounding Run(), so
the lock is held for the process lifetime. The second stream blocks in
resource.lock.LockContext(canceler) and does not return while the Agent runs.
A goroutine profile from elastic-agent diagnostics collect on a three-stream Agent shows exactly
two goroutines parked in go-concert/unison.Mutex.LockContext, called from
input-cursor.lockResource -> (*InputManager).lock -> (*managedInput).runSource: one holder and
two waiters for three configured streams.
4. Why both streams get the same Source.Name(). packages/tenable_io/manifest.yml declares one
policy-template variable:
- name: url
type: text
title: URL
default: https://cloud.tenable.com
Every data stream's agent/stream/cel.yml.hbs renders resource.url: {{url}} from it, and no
data_stream/*/manifest.yml declares a stream-level url to override it -- checked for asset,
vulnerability, audit, plugin and scan. One base URL for one vendor API is an ordinary
package design; it produces identical Source.Name() values for every stream in the package.
Proposed fixes
1. Scope the lock key by data stream (x-pack/filebeat/input/cel/input_manager.go), falling
back to current behaviour when no dataset is known, which keeps deduplication by URL for
hand-configured cel inputs outside Fleet:
func (s *source) Name() string {
if s.cfg.DataStream != "" {
return s.cfg.DataStream + "::" + s.cfg.Resource.URL.String()
}
return s.cfg.Resource.URL.String()
}
s.cfg.DataStream is already populated on every Fleet-managed configure call, so this needs no new
plumbing. It changes exclusivity from one input per URL to one input per (dataset, URL).
Compatibility cost. The persisted cursor-store key (createSourceID:
"<type>::[<userID>::]<Source.Name()>") changes for every affected deployment, so both streams
start from initial_interval once on upgrade rather than resuming. The affected stream has no
cursor to lose; the running one re-collects one initial_interval window, which is the same cost as
a restart without persistence. Whether that needs a version gate is a call for Elastic.
2. Per-package mitigation, without a framework change. Give each data stream a distinct
resource.url by moving the dataset-specific path into the templated URL rather than building it in
the CEL program from a shared base. This changes Source.Name() with no framework change, at the
cost of rewriting each stream's request-building CEL against its now-specific base URL. It is a
non-trivial PR per package, and covers only the packages changed.
Integration Name
Tenable Vulnerability Management [packages/tenable_io]
Dataset Name
All
Integration Version
4.14.0
Agent Version
9.5.3
Agent Output Type
logstash
Elasticsearch Version
9.5.2
OS Version and Architecture
Red Hat Enterprise Linux 9.6
Software/API Version
N/A
Error Message
No error message was displayed, not with debug logging enabled, not with CEL tracing enabled.
Event Original
N/A
There are no problems with the events being generated by the integration.
What did you do?
Applied stock integration code with three streams enabled (asset, vulnerability, plugin). Below you will find the code we are using right now, with a small diagnostic patch added to test the theory that API throttling was causing our problem, we we did see throttling messages from the API server. However, the patch did not fix the issue.
What did you see?
Only one data stream collects (asset). Remaining data streams never collect data until the other data streams are disabled. i.e to collect "vulnerability" (third in the array or streams we enabled), we need to disable both "asset" and "plugin" streams.
What did you expect to see?
Three data streams collecting consistently (asset, plugin, vulnerability)
Anything else?
Yes... we applied Claude Code to diagnose the problem, and on review the agent suggested that the problem is that the API server URL is shared across all three streams, which triggers a bug in the CEL input cursor framework. The evidence presented seems strong, so I will provide the agent diagnostics below. Before I do that I will note that there are two acceptable fixes from my standpoint:
resource.urlby moving the dataset-specific path into the templated URL rather than building it in the CEL program from a shared base. This is the easiest fix as it is specific to this integration, but I think we likely would see the same issue on other multi-stream integrations that share the same URL. I cannot verify this as all the CEL-based integrations that my team has used so far all use a different URL per-stream (see: Abnormal Security, Crowdstrike, Wiz).The following the Claude's own root cause analysis and proposed fix. I lack the golang background to evaluate the accuracy of the report fully, but I can tell you that the AI identified the potential bug in the source first, and then verified though a stack analysis gathered from a running agent using
elastic-agent diagnostics collect, so it's claims feel defensible to me.Root cause
The package's CEL
programis not involved.tenable_io'sassetprogram signals completion of anexport cycle with
"want_more": falseonce its cursor advances, which is correct. The exclusivityis two layers below, in the input-cursor framework.
1. The lock key is the URL alone (
x-pack/filebeat/input/cel/input_manager.go):Source.Name()is what the framework uses to decide whether two configured streams are the sameresource. For
assetandvulnerabilityit returns the same string,https://cloud.tenable.com, for the reason in point 4.config.DataStreamcarries the stream's owndata_stream.dataseton the same struct, populated incursorConfigure(src.cfg.DataStream = dataStreamName(cfg)) and already used elsewhere in theinput's lifecycle for status reporting.
Name()does not read it.2. The lock is held for the input's lifetime (
filebeat/input/v2/input-cursor/):3.
Run()does not return during normal operation(
x-pack/filebeat/input/cel/input.go):timed.Periodicblocks until the input's context ends, re-runningrunCycleeveryinterval. ACEL program finishing its export window ends one
runCyclecall, not the surroundingRun(), sothe lock is held for the process lifetime. The second stream blocks in
resource.lock.LockContext(canceler)and does not return while the Agent runs.A goroutine profile from
elastic-agent diagnostics collecton a three-stream Agent shows exactlytwo goroutines parked in
go-concert/unison.Mutex.LockContext, called frominput-cursor.lockResource->(*InputManager).lock->(*managedInput).runSource: one holder andtwo waiters for three configured streams.
4. Why both streams get the same
Source.Name().packages/tenable_io/manifest.ymldeclares onepolicy-template variable:
Every data stream's
agent/stream/cel.yml.hbsrendersresource.url: {{url}}from it, and nodata_stream/*/manifest.ymldeclares a stream-levelurlto override it -- checked forasset,vulnerability,audit,pluginandscan. One base URL for one vendor API is an ordinarypackage design; it produces identical
Source.Name()values for every stream in the package.Proposed fixes
1. Scope the lock key by data stream (
x-pack/filebeat/input/cel/input_manager.go), fallingback to current behaviour when no dataset is known, which keeps deduplication by URL for
hand-configured
celinputs outside Fleet:s.cfg.DataStreamis already populated on every Fleet-managed configure call, so this needs no newplumbing. It changes exclusivity from one input per URL to one input per (dataset, URL).
Compatibility cost. The persisted cursor-store key (
createSourceID:"<type>::[<userID>::]<Source.Name()>") changes for every affected deployment, so both streamsstart from
initial_intervalonce on upgrade rather than resuming. The affected stream has nocursor to lose; the running one re-collects one
initial_intervalwindow, which is the same cost asa restart without persistence. Whether that needs a version gate is a call for Elastic.
2. Per-package mitigation, without a framework change. Give each data stream a distinct
resource.urlby moving the dataset-specific path into the templated URL rather than building it inthe CEL program from a shared base. This changes
Source.Name()with no framework change, at thecost of rewriting each stream's request-building CEL against its now-specific base URL. It is a
non-trivial PR per package, and covers only the packages changed.