Remove pickle from parallel external source shared memory messages - #6512
devin-ai-integration[bot] wants to merge 2 commits into
Conversation
Replace pickle-based encoding of ScheduledTask, CompletedTask and sample meta-data with a validated, fixed-layout struct-based binary format. Use memfd_create for shared memory allocation, falling back to shm_open. Signed-off-by: Joaquin Anton <janton@nvidia.com> Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
|
I'll fix CI failures and address comments from users with write access. I'll skip comments containing "(aside)".
|
| raise TypeError(f"Unsupported batch argument type: `{type(arg)}`.") | ||
|
|
||
|
|
||
| def _read_batch_args(reader): |
There was a problem hiding this comment.
Intentional: _read_batch_args returns the positional-argument tuple for the batch callback, which is either empty (callback takes no argument) or holds a single iteration / BatchInfo.
|
…tests Signed-off-by: Joaquin Anton <janton@nvidia.com> Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
| if depth > _MAX_SAMPLE_NESTING: | ||
| raise TypeError( | ||
| f"Samples nested deeper than {_MAX_SAMPLE_NESTING} levels are not supported." | ||
| ) |
There was a problem hiding this comment.
Nested samples lose their error
When a parallel callback returns a sample nested more than 64 levels deep, this guard raises in the worker's dispatcher thread, outside the callback-error handler. The dispatcher sends no CompletedTask, so the parent reports a generic worker-interruption error instead of explaining the nesting problem. Catch serialization failures in the dispatcher and return them as failed tasks.
Knowledge Base Used: Python API and framework iterators
| def unpack(self, fmt): | ||
| fmt = struct.Struct("<" + fmt) | ||
| return fmt.unpack(self._take(fmt.size)) | ||
|
|
There was a problem hiding this comment.
Protocol tags use magic integers
The new protocol uses integer constants to dispatch message, task, argument, exception, and sample kinds. This violates the repository directive to use Enum or IntEnum for type and mode dispatch. Use typed tag groups so the protocol's cases remain explicit.
Rule Used: Use enums (or IntEnum/Enum in Python) for type/mode/format dispatch — not magic strings or magic integers. Self-documenting and compiler-checkable. (source)
Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!
Category:
Bug fix (non-breaking change which fixes an issue) — security hardening
Description:
Parallel ExternalSource (PES) exchanged its internal task/result messages and sample meta-data through shared memory encoded with
pickle. Anyone able to modify a shared-memory chunk (e.g. a compromised worker) could make the other side execute arbitrary code onpickle.loads.This PR replaces pickle in the PES shared-memory protocol with a small, fixed-layout, little-endian
struct-based format (_BinaryWriter/_BinaryReaderinshared_batch.py):The reader validates every read against the buffer bounds, rejects unknown tags, trailing bytes, invalid
SampleRanges, invalid / object dtypes and overly deep nesting, raisingRuntimeError("Malformed shared memory message: ..."). Worker exceptions were already normalized toStopIteration/RuntimeErrorbyCompletedTask.failed, so the format carries only those two types plus their message.Additionally (nice-to-have from the ticket),
ShmHandle::CreateHandlenow usesmemfd_createso the shared memory never has a name in the filesystem, even briefly; it falls back to the previousmkstemp+shm_open+shm_unlinkpath whenmemfd_createis unavailable (ENOSYS, orEPERM/EACCES/EINVALfrom seccomp/old kernels). The fd is still passed to workers over Unix sockets, which works the same for memfds.Callback passing reassessment: callbacks (and the optional
py_callback_pickleroutput) are still pickled, but only in the parent -> worker direction, asmultiprocessing.Processarguments at worker start-up (spawn/forkserver) — they never go through the shared-memory chunks and the parent never unpickles anything coming from workers. Executing the user's callback is arbitrary code execution by design, so this is left unchanged.Additional information:
Affected modules and functionalities:
dali/python/nvidia/dali/_multiproc/shared_batch.py: new binary codec (serialize_message/deserialize_message,serialize_samples_meta/deserialize_samples_meta);read_shm_message,write_shm_messageandSharedBatchWriteruse it. The previous restricted unpickler is removed.dali/python/nvidia/dali/_multiproc/messages.py,worker.py: docstrings only.dali/core/os/shared_mem.cc:memfd_createwithshm_openfallback.Key points relevant for the review:
ExternalSourceactually produces: no argument, anintiteration, orBatchInfo.Tests:
ScheduledTask(sample range, sliced range, int /BatchInfo/ empty batch args),CompletedTask(success/failure), nested tuple/list samples and more dtypes; rejection of pickle payloads, truncated/trailing data, bad tags, invalid dtypes, object dtypes, deep nesting.Checklist
Documentation
DALI team only
Requirements
REQ IDs: N/A
JIRA TASK: DALI-4559
Link to Devin session: https://nvidia-cloud.devinenterprise.com/sessions/28dadd725c2143758024982108533e2c
Open in Devin Desktop: https://nvidia-cloud.devinenterprise.com/desktop/session/28dadd725c2143758024982108533e2c?variant=devin