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
83 changes: 83 additions & 0 deletions doc/api/diagnostics_channel.md
Original file line number Diff line number Diff line change
Expand Up @@ -1636,6 +1636,89 @@ diagnosticsChannel.subscribe('crypto.fips.indicator', (message) => {
});
```

#### Filesystem

> Stability: 1 - Experimental

These channels are emitted for file system operations performed through
`node:fs` and `node:fs/promises`. Each operation has its own
[`TracingChannel`][] family named `fs.<operation>`, where `<operation>` is a
stable operation name such as `open`, `read`, `write`, `stat`, `readdir`, or
`realpath`. Subscribers can use [`diagnostics_channel.tracingChannel()`][] to
subscribe to all events of a given operation at once:

```mjs
import diagnostics_channel from 'node:diagnostics_channel';

const channel = diagnostics_channel.tracingChannel('fs.open');
channel.subscribe({
start: (event) => console.log('start', event),
end: (event) => console.log('end', event),
error: (event) => console.log('error', event),
});
```

The events are published from the internal file system implementation, so they
are observed for every public `fs` operation regardless of whether the
function reference was captured before subscribing or whether the operation
uses the callback, promise, or synchronous API.

Each event carries an object with the following common fields:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Other args may be necessary for o11y, for example the args to chmod/chown other than the path.


* `api` {string} The API that performed the operation: `'sync'`, `'callback'`,
or `'promise'`.
* `path` {string|undefined} The path argument for path-based operations, or
the source path for operations with a destination.
* `dest` {string|undefined} The destination argument for operations that
accept one, such as `rename`, `link`, `symlink`, or `copyFile`.
* `fd` {number|undefined} The file descriptor for operations that operate on
an existing file descriptor, such as `read`, `write`, `fsync`, or `close`.

Large read/write buffers are not copied into the event payload. The `start`

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why not include them? Not copy, but pass the reference.

and `asyncStart` events carry no `result` or `error`; the `end` and `asyncEnd`
events carry the `result` of the operation, and the `error` event carries the
`error`, following the [TracingChannel Channels][] conventions.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The convention includes result fields in asyncStart https://nodejs.org/docs/latest/api/diagnostics_channel.html#asyncstartevent

In fact this is necessary for many kinds of instrumentation.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This one is the only blocker, the rest can be handled in future PRs if desired.


Operations performed through streams (`fs.createReadStream` and
`fs.createWriteStream`), most `FileHandle` methods, and the `fs.readFile`

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why not FileHandle?

The streams should be doable with diagnostics channels, but maybe just not TracingChannels.

fast path (which batches open/stat/read/close into a single background job)
are not covered by these channel families, and may not emit the full set of
events.
Comment on lines +1682 to +1686

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This says the readFile fast path isn't covered, but the paragraph above says every public fs operation is.

Also, writeFile and fs.promises.readFile skip the channels, and so do readFileSync with an encoding and existsSync. Those are the fs calls we instrument the most.

I think the PR should cover them too, so all the public fs APIs emit.


##### Event: `'tracing:fs.<operation>:start'`

Emitted synchronously when an operation begins, before the operation is
submitted. For synchronous operations this is followed by `end` (or `error`);
for asynchronous operations it is followed by `end` and then `asyncStart`/
`asyncEnd` (or `error`).

##### Event: `'tracing:fs.<operation>:end'`

* `result` {any} The result of the operation.

Emitted when the operation completes. For synchronous operations this carries
the operation `result`; for asynchronous operations it is emitted when the
operation is submitted and carries no `result` (the `result` is delivered on
the `asyncEnd` event).

##### Event: `'tracing:fs.<operation>:asyncStart'`

Emitted when the asynchronous work for an operation begins (when the
completion callback is invoked).

##### Event: `'tracing:fs.<operation>:asyncEnd'`

* `result` {any} The result of the operation.

Emitted when the asynchronous work for an operation completes, carrying the
operation `result`.

##### Event: `'tracing:fs.<operation>:error'`

* `error` {Error} The error that caused the operation to fail.

Emitted when an operation fails.

#### HTTP

> Stability: 1 - Experimental
Expand Down
1 change: 1 addition & 0 deletions src/env_properties.h
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@
V(allow_bare_named_params_string, "allowBareNamedParameters") \
V(allow_unknown_named_params_string, "allowUnknownNamedParameters") \
V(alpn_callback_string, "ALPNCallback") \
V(api_string, "api") \
V(args_string, "args") \
V(arguments_string, "arguments") \
V(async_ids_stack_string, "async_ids_stack") \
Expand Down
132 changes: 132 additions & 0 deletions src/node_file-inl.h
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@
#include "node_file.h"
#include "req_wrap-inl.h"

#include <type_traits>

namespace node {
namespace fs {

Expand Down Expand Up @@ -204,6 +206,37 @@ FSReqPromise<AliasedBufferT>::~FSReqPromise() {
CHECK_IMPLIES(!finished_, !env()->can_call_into_js());
}

inline bool FSOperationChannelHasSubscribers(FSOperationChannels& channels,
FSOperationChannel channel) {
diagnostics_channel::Channel* ch =
channels[static_cast<size_t>(channel)].get();
return ch != nullptr && ch->HasSubscribers();
}

inline bool AnyFSOperationChannelHasSubscribers(FSOperationChannels& channels) {
for (size_t i = 0; i < kNumFSOperationChannels; i++) {
diagnostics_channel::Channel* ch = channels[i].get();
if (ch != nullptr && ch->HasSubscribers()) return true;
}
return false;
}

// Returns the file descriptor an fs operation acts on, or -1 for operations
// that take a path. libuv takes the descriptor as the first argument after
// the request for every descriptor-based operation (close, read, write,
// fstat, ...), while path-based operations take a string there.
inline int FSOperationFd() {
return -1;
}
template <typename First, typename... Rest>
inline int FSOperationFd(First first, Rest...) {
if constexpr (std::is_same_v<First, int>) {
return first;
} else {
return -1;
}
}

template <typename AliasedBufferT>
FSReqPromise<AliasedBufferT>::FSReqPromise(BindingData* binding_data,
v8::Local<v8::Object> obj,
Expand All @@ -214,6 +247,7 @@ FSReqPromise<AliasedBufferT>::FSReqPromise(BindingData* binding_data,
template <typename AliasedBufferT>
void FSReqPromise<AliasedBufferT>::Reject(v8::Local<v8::Value> reject) {
finished_ = true;
PublishFSOpCompletionEvent(this, FSOperationChannel::kError, "error", reject);
v8::HandleScope scope(env()->isolate());
InternalCallbackScope callback_scope(this);
v8::Local<v8::Value> value;
Expand All @@ -232,6 +266,8 @@ void FSReqPromise<AliasedBufferT>::Reject(v8::Local<v8::Value> reject) {
template <typename AliasedBufferT>
void FSReqPromise<AliasedBufferT>::Resolve(v8::Local<v8::Value> value) {
finished_ = true;
PublishFSOpCompletionEvent(
this, FSOperationChannel::kAsyncEnd, "result", value);
v8::HandleScope scope(env()->isolate());
InternalCallbackScope callback_scope(this);
v8::Local<v8::Value> val;
Expand Down Expand Up @@ -311,6 +347,7 @@ FSReqBase* GetReqWrap(const v8::FunctionCallbackInfo<v8::Value>& args,
result =
FSReqPromise<AliasedFloat64Array>::New(binding_data, use_bigint);
}
result->set_is_promise(true);
}
}
if (result != nullptr) {
Expand All @@ -328,13 +365,51 @@ FSReqBase* AsyncDestCall(Environment* env, FSReqBase* req_wrap,
Func fn, Args... fn_args) {
CHECK_NOT_NULL(req_wrap);
req_wrap->Init(syscall, dest, len, enc);
BindingData* binding = req_wrap->binding_data();
const char* api = req_wrap->is_promise() ? "promise" : "callback";
FSOperationChannels* channels = nullptr;
// See SyncCallAndThrowIf: instrumentation is unsafe with a pending
// exception.
if (binding != nullptr && !env->isolate()->HasPendingException()) {
channels = GetFSOperationChannels(binding, env, syscall);
req_wrap->set_op_channels(channels);
if (channels != nullptr && FSOperationChannelHasSubscribers(
*channels, FSOperationChannel::kStart)) {
PublishFSOperationEvent(env,
*channels,
FSOperationChannel::kStart,
api,
nullptr,
req_wrap->data(),
-1,
nullptr,
v8::Local<v8::Value>());
}
}
int err = req_wrap->Dispatch(fn, fn_args..., after);
if (err < 0) {
uv_fs_t* uv_req = req_wrap->req();
uv_req->result = err;
uv_req->path = nullptr;
after(uv_req); // after may delete req_wrap if there is an error
req_wrap = nullptr;
} else if (channels != nullptr &&
AnyFSOperationChannelHasSubscribers(*channels)) {
const char* path = req_wrap->req()->path;
int fd = FSOperationFd(fn_args...);
req_wrap->set_fd(fd);
// The path is captured for the completion events; it requires a copy
// since the uv request is cleaned up before they fire.
req_wrap->set_op_path(path == nullptr ? std::string() : path);
PublishFSOperationEvent(env,
*channels,
FSOperationChannel::kEnd,
api,
path,
req_wrap->data(),
fd,
nullptr,
v8::Local<v8::Value>());
}
return req_wrap;
}
Expand Down Expand Up @@ -389,7 +464,64 @@ int SyncCallAndThrowIf(Predicate should_throw,
Func fn,
Args... args) {
env->PrintSyncTrace();
BindingData* binding = Realm::GetBindingData<BindingData>(env->context());
FSOperationChannels* channels = nullptr;
// The instrumentation creates V8 objects and may run subscribers, neither
// of which is safe with a pending exception (a multi-step operation keeps
// going after a failed step to clean up, e.g. write + close).
if (binding != nullptr && !env->isolate()->HasPendingException()) {
channels = GetFSOperationChannels(binding, env, req_wrap->syscall_p);
if (channels != nullptr && FSOperationChannelHasSubscribers(
*channels, FSOperationChannel::kStart)) {
PublishFSOperationEvent(env,
*channels,
FSOperationChannel::kStart,
"sync",
req_wrap->path_p,
req_wrap->dest_p,
-1,
nullptr,
v8::Local<v8::Value>());
}
}
int result = fn(nullptr, &(req_wrap->req), args..., nullptr);
if (channels != nullptr) {
if (should_throw(result)) {
// The error object is only built when someone is listening; the throw
// path below creates its own copy.
if (FSOperationChannelHasSubscribers(*channels,
FSOperationChannel::kError)) {
int fd = FSOperationFd(args...);
v8::Local<v8::Value> error = UVException(env->isolate(),
result,
req_wrap->syscall_p,
nullptr,
req_wrap->path_p,
req_wrap->dest_p);
PublishFSOperationEvent(env,
*channels,
FSOperationChannel::kError,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

On failure this emits error and then stops, with no end after it. traceSync always emits end after error, for APMs end is important because that's when they end the span or flush any telemetry.

"sync",
req_wrap->path_p,
req_wrap->dest_p,
fd,
"error",
error);
}
} else if (FSOperationChannelHasSubscribers(*channels,
FSOperationChannel::kEnd)) {
int fd = FSOperationFd(args...);
PublishFSOperationEvent(env,
*channels,
FSOperationChannel::kEnd,
"sync",
req_wrap->path_p,
req_wrap->dest_p,
fd,
"result",
v8::Integer::New(env->isolate(), result));
}
}
if (should_throw(result)) {
env->ThrowUVException(result,
req_wrap->syscall_p,
Expand Down
Loading
Loading