[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe - #16992

Merged
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer
May 14, 2024
Merged

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe#16992
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer

Conversation

@Lunderberg

Copy link
Copy Markdown
Contributor

Prior to this commit, using disco.Session methods to transfer NDArray instances to workers could raise an exception if the NDArray is larger than the buffer allocated by the OS for the controller/worker pipe. In these case, the first call to the Read method of tvm::support::Pipe would successfully return, but only with the initial bytes of the NDArray. Receiving the full NDArray requires repeatedly calling the POSIX read function.

This commit updates the Read and Write methods of tvm::support::Pipe to repeatedly call the underlying read/write methods, until the full NDArray has been transferred.

This commit does not add any unit tests, as the existing unit test tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession] requires this change to pass.

Prior to this commit, using `disco.Session` methods to transfer
`NDArray` instances to workers could raise an exception if the
`NDArray` is larger than the buffer allocated by the OS for the
controller/worker pipe. In these case, the first call to the `Read`
method of `tvm::support::Pipe` would successfully return, but only
with the initial bytes of the `NDArray`. Receiving the full `NDArray`
requires repeatedly calling the POSIX `read` function.
This commit updates the `Read` and `Write` methods of
`tvm::support::Pipe` to repeatedly call the underlying read/write
methods, until the full `NDArray` has been transferred.
This commit does not add any unit tests, as the existing unit test
`tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession]`
requires this change to pass.
@Lunderberg
Lunderberg merged commit d9dbbc9 into apache:mainMay 14, 2024
@Lunderberg
Lunderberg deleted the bugfix_disco_transfer_larger_than_pipe_buffer branch May 14, 2024 14:39
@tqchen

Copy link
Copy Markdown
Member

Thanks @Lunderberg !
In this case, i think it is better to introduce a ReadAll and WriteAll function, in pairing with https://github.com/apache/tvm/blob/main/src/support/socket.h#L492, and then we call these functions instead

Just to keep Socket and pipe behavior consistent with each other, low-level read/write can remain partial while readall and writeall contains convenient method for the intended behavior

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

I agree. This is a stop-gap measure, as the longer-term fix will require updating the dmlc::Stream API to return size_t from Write, similar to Read. Without that API change, the calling scope cannot correctly implement a WriteAll method.

Lunderberg added a commit to Lunderberg/tvm that referenced this pull request May 14, 2024
This resolves a conflict between two recent changes. In
apache#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
apache#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

@tqchen The changes for the long-term fix are implemented in dmlc/dmlc-core#686, with TVM compatibility changes implemented in #16998. Once those changes land, we will be able to remove the stop-gap implementation, and provide ReadAll and WriteAll helper methods.

@tqchen

Copy link
Copy Markdown
Member

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes). After seeing your fix, i realized that this PR actually did the right thing given pipe inheritated from the stream interface

@tqchen

Copy link
Copy Markdown
Member

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes).

Thank you, and that history makes a lot of sense. I'm thinking it would still be a good change (though low priority) to expose partial writes in dmlc::Stream interface, because the Stream interface is also used for pipes and TCP sockets.

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

Good point. Looking into it, it looks like the Socket interface also includes functionality that isn't appropriate for a pipe, such as the address/port to listen on, accepting new connections, etc.

Lunderberg added a commit that referenced this pull request May 15, 2024
This resolves a conflict between two recent changes. In
#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Lunderberg@tqchen@masahi
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all \u003cpre\u003e\u003ccode\u003e blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks"); } } catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); } })(); (function(){ try { var __m = "github.com"; var __re = new RegExp('^' + "github\\.com" + '
Skip to content

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe - #16992

Merged
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer
May 14, 2024
Merged

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe#16992
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer

Conversation

@Lunderberg

Copy link
Copy Markdown
Contributor

Prior to this commit, using disco.Session methods to transfer NDArray instances to workers could raise an exception if the NDArray is larger than the buffer allocated by the OS for the controller/worker pipe. In these case, the first call to the Read method of tvm::support::Pipe would successfully return, but only with the initial bytes of the NDArray. Receiving the full NDArray requires repeatedly calling the POSIX read function.

This commit updates the Read and Write methods of tvm::support::Pipe to repeatedly call the underlying read/write methods, until the full NDArray has been transferred.

This commit does not add any unit tests, as the existing unit test tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession] requires this change to pass.

Prior to this commit, using `disco.Session` methods to transfer
`NDArray` instances to workers could raise an exception if the
`NDArray` is larger than the buffer allocated by the OS for the
controller/worker pipe. In these case, the first call to the `Read`
method of `tvm::support::Pipe` would successfully return, but only
with the initial bytes of the `NDArray`. Receiving the full `NDArray`
requires repeatedly calling the POSIX `read` function.
This commit updates the `Read` and `Write` methods of
`tvm::support::Pipe` to repeatedly call the underlying read/write
methods, until the full `NDArray` has been transferred.
This commit does not add any unit tests, as the existing unit test
`tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession]`
requires this change to pass.
@Lunderberg
Lunderberg merged commit d9dbbc9 into apache:mainMay 14, 2024
@Lunderberg
Lunderberg deleted the bugfix_disco_transfer_larger_than_pipe_buffer branch May 14, 2024 14:39
@tqchen

Copy link
Copy Markdown
Member

Thanks @Lunderberg !
In this case, i think it is better to introduce a ReadAll and WriteAll function, in pairing with https://github.com/apache/tvm/blob/main/src/support/socket.h#L492, and then we call these functions instead

Just to keep Socket and pipe behavior consistent with each other, low-level read/write can remain partial while readall and writeall contains convenient method for the intended behavior

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

I agree. This is a stop-gap measure, as the longer-term fix will require updating the dmlc::Stream API to return size_t from Write, similar to Read. Without that API change, the calling scope cannot correctly implement a WriteAll method.

Lunderberg added a commit to Lunderberg/tvm that referenced this pull request May 14, 2024
This resolves a conflict between two recent changes. In
apache#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
apache#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

@tqchen The changes for the long-term fix are implemented in dmlc/dmlc-core#686, with TVM compatibility changes implemented in #16998. Once those changes land, we will be able to remove the stop-gap implementation, and provide ReadAll and WriteAll helper methods.

@tqchen

Copy link
Copy Markdown
Member

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes). After seeing your fix, i realized that this PR actually did the right thing given pipe inheritated from the stream interface

@tqchen

Copy link
Copy Markdown
Member

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes).

Thank you, and that history makes a lot of sense. I'm thinking it would still be a good change (though low priority) to expose partial writes in dmlc::Stream interface, because the Stream interface is also used for pipes and TCP sockets.

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

Good point. Looking into it, it looks like the Socket interface also includes functionality that isn't appropriate for a pipe, such as the address/port to listen on, accepting new connections, etc.

Lunderberg added a commit that referenced this pull request May 15, 2024
This resolves a conflict between two recent changes. In
#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Lunderberg@tqchen@masahi
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe - #16992

Merged
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer
May 14, 2024
Merged

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe#16992
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer

Conversation

@Lunderberg

Copy link
Copy Markdown
Contributor

Prior to this commit, using disco.Session methods to transfer NDArray instances to workers could raise an exception if the NDArray is larger than the buffer allocated by the OS for the controller/worker pipe. In these case, the first call to the Read method of tvm::support::Pipe would successfully return, but only with the initial bytes of the NDArray. Receiving the full NDArray requires repeatedly calling the POSIX read function.

This commit updates the Read and Write methods of tvm::support::Pipe to repeatedly call the underlying read/write methods, until the full NDArray has been transferred.

This commit does not add any unit tests, as the existing unit test tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession] requires this change to pass.

Prior to this commit, using `disco.Session` methods to transfer
`NDArray` instances to workers could raise an exception if the
`NDArray` is larger than the buffer allocated by the OS for the
controller/worker pipe. In these case, the first call to the `Read`
method of `tvm::support::Pipe` would successfully return, but only
with the initial bytes of the `NDArray`. Receiving the full `NDArray`
requires repeatedly calling the POSIX `read` function.
This commit updates the `Read` and `Write` methods of
`tvm::support::Pipe` to repeatedly call the underlying read/write
methods, until the full `NDArray` has been transferred.
This commit does not add any unit tests, as the existing unit test
`tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession]`
requires this change to pass.
@Lunderberg
Lunderberg merged commit d9dbbc9 into apache:mainMay 14, 2024
@Lunderberg
Lunderberg deleted the bugfix_disco_transfer_larger_than_pipe_buffer branch May 14, 2024 14:39
@tqchen

Copy link
Copy Markdown
Member

Thanks @Lunderberg !
In this case, i think it is better to introduce a ReadAll and WriteAll function, in pairing with https://github.com/apache/tvm/blob/main/src/support/socket.h#L492, and then we call these functions instead

Just to keep Socket and pipe behavior consistent with each other, low-level read/write can remain partial while readall and writeall contains convenient method for the intended behavior

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

I agree. This is a stop-gap measure, as the longer-term fix will require updating the dmlc::Stream API to return size_t from Write, similar to Read. Without that API change, the calling scope cannot correctly implement a WriteAll method.

Lunderberg added a commit to Lunderberg/tvm that referenced this pull request May 14, 2024
This resolves a conflict between two recent changes. In
apache#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
apache#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

@tqchen The changes for the long-term fix are implemented in dmlc/dmlc-core#686, with TVM compatibility changes implemented in #16998. Once those changes land, we will be able to remove the stop-gap implementation, and provide ReadAll and WriteAll helper methods.

@tqchen

Copy link
Copy Markdown
Member

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes). After seeing your fix, i realized that this PR actually did the right thing given pipe inheritated from the stream interface

@tqchen

Copy link
Copy Markdown
Member

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes).

Thank you, and that history makes a lot of sense. I'm thinking it would still be a good change (though low priority) to expose partial writes in dmlc::Stream interface, because the Stream interface is also used for pipes and TCP sockets.

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

Good point. Looking into it, it looks like the Socket interface also includes functionality that isn't appropriate for a pipe, such as the address/port to listen on, accepting new connections, etc.

Lunderberg added a commit that referenced this pull request May 15, 2024
This resolves a conflict between two recent changes. In
#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Lunderberg@tqchen@masahi
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length \u003e 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe - #16992

Merged
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer
May 14, 2024
Merged

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe#16992
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer

Conversation

@Lunderberg

Copy link
Copy Markdown
Contributor

Prior to this commit, using disco.Session methods to transfer NDArray instances to workers could raise an exception if the NDArray is larger than the buffer allocated by the OS for the controller/worker pipe. In these case, the first call to the Read method of tvm::support::Pipe would successfully return, but only with the initial bytes of the NDArray. Receiving the full NDArray requires repeatedly calling the POSIX read function.

This commit updates the Read and Write methods of tvm::support::Pipe to repeatedly call the underlying read/write methods, until the full NDArray has been transferred.

This commit does not add any unit tests, as the existing unit test tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession] requires this change to pass.

Prior to this commit, using `disco.Session` methods to transfer
`NDArray` instances to workers could raise an exception if the
`NDArray` is larger than the buffer allocated by the OS for the
controller/worker pipe. In these case, the first call to the `Read`
method of `tvm::support::Pipe` would successfully return, but only
with the initial bytes of the `NDArray`. Receiving the full `NDArray`
requires repeatedly calling the POSIX `read` function.
This commit updates the `Read` and `Write` methods of
`tvm::support::Pipe` to repeatedly call the underlying read/write
methods, until the full `NDArray` has been transferred.
This commit does not add any unit tests, as the existing unit test
`tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession]`
requires this change to pass.
@Lunderberg
Lunderberg merged commit d9dbbc9 into apache:mainMay 14, 2024
@Lunderberg
Lunderberg deleted the bugfix_disco_transfer_larger_than_pipe_buffer branch May 14, 2024 14:39
@tqchen

Copy link
Copy Markdown
Member

Thanks @Lunderberg !
In this case, i think it is better to introduce a ReadAll and WriteAll function, in pairing with https://github.com/apache/tvm/blob/main/src/support/socket.h#L492, and then we call these functions instead

Just to keep Socket and pipe behavior consistent with each other, low-level read/write can remain partial while readall and writeall contains convenient method for the intended behavior

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

I agree. This is a stop-gap measure, as the longer-term fix will require updating the dmlc::Stream API to return size_t from Write, similar to Read. Without that API change, the calling scope cannot correctly implement a WriteAll method.

Lunderberg added a commit to Lunderberg/tvm that referenced this pull request May 14, 2024
This resolves a conflict between two recent changes. In
apache#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
apache#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

@tqchen The changes for the long-term fix are implemented in dmlc/dmlc-core#686, with TVM compatibility changes implemented in #16998. Once those changes land, we will be able to remove the stop-gap implementation, and provide ReadAll and WriteAll helper methods.

@tqchen

Copy link
Copy Markdown
Member

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes). After seeing your fix, i realized that this PR actually did the right thing given pipe inheritated from the stream interface

@tqchen

Copy link
Copy Markdown
Member

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes).

Thank you, and that history makes a lot of sense. I'm thinking it would still be a good change (though low priority) to expose partial writes in dmlc::Stream interface, because the Stream interface is also used for pipes and TCP sockets.

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

Good point. Looking into it, it looks like the Socket interface also includes functionality that isn't appropriate for a pipe, such as the address/port to listen on, accepting new connections, etc.

Lunderberg added a commit that referenced this pull request May 15, 2024
This resolves a conflict between two recent changes. In
#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Lunderberg@tqchen@masahi
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe - #16992

Merged
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer
May 14, 2024
Merged

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe#16992
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer

Conversation

@Lunderberg

Copy link
Copy Markdown
Contributor

Prior to this commit, using disco.Session methods to transfer NDArray instances to workers could raise an exception if the NDArray is larger than the buffer allocated by the OS for the controller/worker pipe. In these case, the first call to the Read method of tvm::support::Pipe would successfully return, but only with the initial bytes of the NDArray. Receiving the full NDArray requires repeatedly calling the POSIX read function.

This commit updates the Read and Write methods of tvm::support::Pipe to repeatedly call the underlying read/write methods, until the full NDArray has been transferred.

This commit does not add any unit tests, as the existing unit test tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession] requires this change to pass.

Prior to this commit, using `disco.Session` methods to transfer
`NDArray` instances to workers could raise an exception if the
`NDArray` is larger than the buffer allocated by the OS for the
controller/worker pipe. In these case, the first call to the `Read`
method of `tvm::support::Pipe` would successfully return, but only
with the initial bytes of the `NDArray`. Receiving the full `NDArray`
requires repeatedly calling the POSIX `read` function.
This commit updates the `Read` and `Write` methods of
`tvm::support::Pipe` to repeatedly call the underlying read/write
methods, until the full `NDArray` has been transferred.
This commit does not add any unit tests, as the existing unit test
`tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession]`
requires this change to pass.
@Lunderberg
Lunderberg merged commit d9dbbc9 into apache:mainMay 14, 2024
@Lunderberg
Lunderberg deleted the bugfix_disco_transfer_larger_than_pipe_buffer branch May 14, 2024 14:39
@tqchen

Copy link
Copy Markdown
Member

Thanks @Lunderberg !
In this case, i think it is better to introduce a ReadAll and WriteAll function, in pairing with https://github.com/apache/tvm/blob/main/src/support/socket.h#L492, and then we call these functions instead

Just to keep Socket and pipe behavior consistent with each other, low-level read/write can remain partial while readall and writeall contains convenient method for the intended behavior

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

I agree. This is a stop-gap measure, as the longer-term fix will require updating the dmlc::Stream API to return size_t from Write, similar to Read. Without that API change, the calling scope cannot correctly implement a WriteAll method.

Lunderberg added a commit to Lunderberg/tvm that referenced this pull request May 14, 2024
This resolves a conflict between two recent changes. In
apache#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
apache#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

@tqchen The changes for the long-term fix are implemented in dmlc/dmlc-core#686, with TVM compatibility changes implemented in #16998. Once those changes land, we will be able to remove the stop-gap implementation, and provide ReadAll and WriteAll helper methods.

@tqchen

Copy link
Copy Markdown
Member

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes). After seeing your fix, i realized that this PR actually did the right thing given pipe inheritated from the stream interface

@tqchen

Copy link
Copy Markdown
Member

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes).

Thank you, and that history makes a lot of sense. I'm thinking it would still be a good change (though low priority) to expose partial writes in dmlc::Stream interface, because the Stream interface is also used for pipes and TCP sockets.

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

Good point. Looking into it, it looks like the Socket interface also includes functionality that isn't appropriate for a pipe, such as the address/port to listen on, accepting new connections, etc.

Lunderberg added a commit that referenced this pull request May 15, 2024
This resolves a conflict between two recent changes. In
#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Lunderberg@tqchen@masahi
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe - #16992

Merged
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer
May 14, 2024
Merged

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe#16992
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer

Conversation

@Lunderberg

Copy link
Copy Markdown
Contributor

Prior to this commit, using disco.Session methods to transfer NDArray instances to workers could raise an exception if the NDArray is larger than the buffer allocated by the OS for the controller/worker pipe. In these case, the first call to the Read method of tvm::support::Pipe would successfully return, but only with the initial bytes of the NDArray. Receiving the full NDArray requires repeatedly calling the POSIX read function.

This commit updates the Read and Write methods of tvm::support::Pipe to repeatedly call the underlying read/write methods, until the full NDArray has been transferred.

This commit does not add any unit tests, as the existing unit test tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession] requires this change to pass.

Prior to this commit, using `disco.Session` methods to transfer
`NDArray` instances to workers could raise an exception if the
`NDArray` is larger than the buffer allocated by the OS for the
controller/worker pipe. In these case, the first call to the `Read`
method of `tvm::support::Pipe` would successfully return, but only
with the initial bytes of the `NDArray`. Receiving the full `NDArray`
requires repeatedly calling the POSIX `read` function.
This commit updates the `Read` and `Write` methods of
`tvm::support::Pipe` to repeatedly call the underlying read/write
methods, until the full `NDArray` has been transferred.
This commit does not add any unit tests, as the existing unit test
`tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession]`
requires this change to pass.
@Lunderberg
Lunderberg merged commit d9dbbc9 into apache:mainMay 14, 2024
@Lunderberg
Lunderberg deleted the bugfix_disco_transfer_larger_than_pipe_buffer branch May 14, 2024 14:39
@tqchen

Copy link
Copy Markdown
Member

Thanks @Lunderberg !
In this case, i think it is better to introduce a ReadAll and WriteAll function, in pairing with https://github.com/apache/tvm/blob/main/src/support/socket.h#L492, and then we call these functions instead

Just to keep Socket and pipe behavior consistent with each other, low-level read/write can remain partial while readall and writeall contains convenient method for the intended behavior

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

I agree. This is a stop-gap measure, as the longer-term fix will require updating the dmlc::Stream API to return size_t from Write, similar to Read. Without that API change, the calling scope cannot correctly implement a WriteAll method.

Lunderberg added a commit to Lunderberg/tvm that referenced this pull request May 14, 2024
This resolves a conflict between two recent changes. In
apache#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
apache#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

@tqchen The changes for the long-term fix are implemented in dmlc/dmlc-core#686, with TVM compatibility changes implemented in #16998. Once those changes land, we will be able to remove the stop-gap implementation, and provide ReadAll and WriteAll helper methods.

@tqchen

Copy link
Copy Markdown
Member

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes). After seeing your fix, i realized that this PR actually did the right thing given pipe inheritated from the stream interface

@tqchen

Copy link
Copy Markdown
Member

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes).

Thank you, and that history makes a lot of sense. I'm thinking it would still be a good change (though low priority) to expose partial writes in dmlc::Stream interface, because the Stream interface is also used for pipes and TCP sockets.

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

Good point. Looking into it, it looks like the Socket interface also includes functionality that isn't appropriate for a pipe, such as the address/port to listen on, accepting new connections, etc.

Lunderberg added a commit that referenced this pull request May 15, 2024
This resolves a conflict between two recent changes. In
#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Lunderberg@tqchen@masahi
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe - #16992

Merged
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer
May 14, 2024
Merged

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe#16992
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer

Conversation

@Lunderberg

Copy link
Copy Markdown
Contributor

Prior to this commit, using disco.Session methods to transfer NDArray instances to workers could raise an exception if the NDArray is larger than the buffer allocated by the OS for the controller/worker pipe. In these case, the first call to the Read method of tvm::support::Pipe would successfully return, but only with the initial bytes of the NDArray. Receiving the full NDArray requires repeatedly calling the POSIX read function.

This commit updates the Read and Write methods of tvm::support::Pipe to repeatedly call the underlying read/write methods, until the full NDArray has been transferred.

This commit does not add any unit tests, as the existing unit test tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession] requires this change to pass.

Prior to this commit, using `disco.Session` methods to transfer
`NDArray` instances to workers could raise an exception if the
`NDArray` is larger than the buffer allocated by the OS for the
controller/worker pipe. In these case, the first call to the `Read`
method of `tvm::support::Pipe` would successfully return, but only
with the initial bytes of the `NDArray`. Receiving the full `NDArray`
requires repeatedly calling the POSIX `read` function.
This commit updates the `Read` and `Write` methods of
`tvm::support::Pipe` to repeatedly call the underlying read/write
methods, until the full `NDArray` has been transferred.
This commit does not add any unit tests, as the existing unit test
`tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession]`
requires this change to pass.
@Lunderberg
Lunderberg merged commit d9dbbc9 into apache:mainMay 14, 2024
@Lunderberg
Lunderberg deleted the bugfix_disco_transfer_larger_than_pipe_buffer branch May 14, 2024 14:39
@tqchen

Copy link
Copy Markdown
Member

Thanks @Lunderberg !
In this case, i think it is better to introduce a ReadAll and WriteAll function, in pairing with https://github.com/apache/tvm/blob/main/src/support/socket.h#L492, and then we call these functions instead

Just to keep Socket and pipe behavior consistent with each other, low-level read/write can remain partial while readall and writeall contains convenient method for the intended behavior

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

I agree. This is a stop-gap measure, as the longer-term fix will require updating the dmlc::Stream API to return size_t from Write, similar to Read. Without that API change, the calling scope cannot correctly implement a WriteAll method.

Lunderberg added a commit to Lunderberg/tvm that referenced this pull request May 14, 2024
This resolves a conflict between two recent changes. In
apache#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
apache#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

@tqchen The changes for the long-term fix are implemented in dmlc/dmlc-core#686, with TVM compatibility changes implemented in #16998. Once those changes land, we will be able to remove the stop-gap implementation, and provide ReadAll and WriteAll helper methods.

@tqchen

Copy link
Copy Markdown
Member

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes). After seeing your fix, i realized that this PR actually did the right thing given pipe inheritated from the stream interface

@tqchen

Copy link
Copy Markdown
Member

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes).

Thank you, and that history makes a lot of sense. I'm thinking it would still be a good change (though low priority) to expose partial writes in dmlc::Stream interface, because the Stream interface is also used for pipes and TCP sockets.

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

Good point. Looking into it, it looks like the Socket interface also includes functionality that isn't appropriate for a pipe, such as the address/port to listen on, accepting new connections, etc.

Lunderberg added a commit that referenced this pull request May 15, 2024
This resolves a conflict between two recent changes. In
#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Lunderberg@tqchen@masahi
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe - #16992

Merged
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer
May 14, 2024
Merged

[Bugfix][Disco] Handle NDArray larger than OS buffer for pipe#16992
Lunderberg merged 1 commit into
apache:mainfrom
Lunderberg:bugfix_disco_transfer_larger_than_pipe_buffer

Conversation

@Lunderberg

Copy link
Copy Markdown
Contributor

Prior to this commit, using disco.Session methods to transfer NDArray instances to workers could raise an exception if the NDArray is larger than the buffer allocated by the OS for the controller/worker pipe. In these case, the first call to the Read method of tvm::support::Pipe would successfully return, but only with the initial bytes of the NDArray. Receiving the full NDArray requires repeatedly calling the POSIX read function.

This commit updates the Read and Write methods of tvm::support::Pipe to repeatedly call the underlying read/write methods, until the full NDArray has been transferred.

This commit does not add any unit tests, as the existing unit test tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession] requires this change to pass.

Prior to this commit, using `disco.Session` methods to transfer
`NDArray` instances to workers could raise an exception if the
`NDArray` is larger than the buffer allocated by the OS for the
controller/worker pipe. In these case, the first call to the `Read`
method of `tvm::support::Pipe` would successfully return, but only
with the initial bytes of the `NDArray`. Receiving the full `NDArray`
requires repeatedly calling the POSIX `read` function.
This commit updates the `Read` and `Write` methods of
`tvm::support::Pipe` to repeatedly call the underlying read/write
methods, until the full `NDArray` has been transferred.
This commit does not add any unit tests, as the existing unit test
`tests/python/disco/test_ccl.py::test_attention[nccl-ProcessSession]`
requires this change to pass.
@Lunderberg
Lunderberg merged commit d9dbbc9 into apache:mainMay 14, 2024
@Lunderberg
Lunderberg deleted the bugfix_disco_transfer_larger_than_pipe_buffer branch May 14, 2024 14:39
@tqchen

Copy link
Copy Markdown
Member

Thanks @Lunderberg !
In this case, i think it is better to introduce a ReadAll and WriteAll function, in pairing with https://github.com/apache/tvm/blob/main/src/support/socket.h#L492, and then we call these functions instead

Just to keep Socket and pipe behavior consistent with each other, low-level read/write can remain partial while readall and writeall contains convenient method for the intended behavior

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

I agree. This is a stop-gap measure, as the longer-term fix will require updating the dmlc::Stream API to return size_t from Write, similar to Read. Without that API change, the calling scope cannot correctly implement a WriteAll method.

Lunderberg added a commit to Lunderberg/tvm that referenced this pull request May 14, 2024
This resolves a conflict between two recent changes. In
apache#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
apache#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

@tqchen The changes for the long-term fix are implemented in dmlc/dmlc-core#686, with TVM compatibility changes implemented in #16998. Once those changes land, we will be able to remove the stop-gap implementation, and provide ReadAll and WriteAll helper methods.

@tqchen

Copy link
Copy Markdown
Member

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes). After seeing your fix, i realized that this PR actually did the right thing given pipe inheritated from the stream interface

@tqchen

Copy link
Copy Markdown
Member

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

@Lunderberg

Copy link
Copy Markdown
ContributorAuthor

Ah I see, the Stream interface was actually created to be in analogy of FileStream interface(which do not have partial writes).

Thank you, and that history makes a lot of sense. I'm thinking it would still be a good change (though low priority) to expose partial writes in dmlc::Stream interface, because the Stream interface is also used for pipes and TCP sockets.

The socket interface was mainly meant to accomodate for the non-blocking case where full read is not always possible

Good point. Looking into it, it looks like the Socket interface also includes functionality that isn't appropriate for a pipe, such as the address/port to listen on, accepting new connections, etc.

Lunderberg added a commit that referenced this pull request May 15, 2024
This resolves a conflict between two recent changes. In
#16989, reads of size zero are used
to identify hangups in `ProcessSession`. In
#16992, reads of size zero are
treated as an error to avoid infinite loops while waiting for data to
be ready.
For a long-term resolution, the `dmlc::Stream` interface will need to
be updated, so that the `Write` method returns the number of bytes
written, just as the `Read` method currently does. This will allow
the calling scope to verify the number of bytes received.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@Lunderberg@tqchen@masahi