Add lock on tee peer cleanup (#9446)

<!-- Thank you for contributing to LangChain!

Replace this entire comment with:
  - Description: a description of the change, 
  - Issue: the issue # it fixes (if applicable),
  - Dependencies: any dependencies required for this change,
- Tag maintainer: for a quicker response, tag the relevant maintainer
(see below),
- Twitter handle: we announce bigger features on Twitter. If your PR
gets announced and you'd like a mention, we'll gladly shout you out!

Please make sure your PR is passing linting and testing before
submitting. Run `make format`, `make lint` and `make test` to check this
locally.

See contribution guidelines for more information on how to write/run
tests, lint, etc:

https://github.com/hwchase17/langchain/blob/master/.github/CONTRIBUTING.md

If you're adding a new integration, please include:
1. a test for the integration, preferably unit tests that do not rely on
network access,
2. an example notebook showing its use. These live is docs/extras
directory.

If no one reviews your PR within a few days, please @-mention one of
@baskaryan, @eyurtsev, @hwchase17, @rlancemartin.
 -->
This commit is contained in:
Nuno Campos 2023-08-18 14:20:09 +01:00 committed by GitHub
parent 0689628489
commit 9417961b17
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
2 changed files with 18 additions and 16 deletions

View File

@ -107,14 +107,15 @@ async def tee_peer(
peer_buffer.append(item) peer_buffer.append(item)
yield buffer.popleft() yield buffer.popleft()
finally: finally:
# this peer is done remove its buffer async with lock:
for idx, peer_buffer in enumerate(peers): # pragma: no branch # this peer is done remove its buffer
if peer_buffer is buffer: for idx, peer_buffer in enumerate(peers): # pragma: no branch
peers.pop(idx) if peer_buffer is buffer:
break peers.pop(idx)
# if we are the last peer, try and close the iterator break
if not peers and hasattr(iterator, "aclose"): # if we are the last peer, try and close the iterator
await iterator.aclose() if not peers and hasattr(iterator, "aclose"):
await iterator.aclose()
class Tee(Generic[T]): class Tee(Generic[T]):

View File

@ -60,14 +60,15 @@ def tee_peer(
peer_buffer.append(item) peer_buffer.append(item)
yield buffer.popleft() yield buffer.popleft()
finally: finally:
# this peer is done remove its buffer with lock:
for idx, peer_buffer in enumerate(peers): # pragma: no branch # this peer is done remove its buffer
if peer_buffer is buffer: for idx, peer_buffer in enumerate(peers): # pragma: no branch
peers.pop(idx) if peer_buffer is buffer:
break peers.pop(idx)
# if we are the last peer, try and close the iterator break
if not peers and hasattr(iterator, "close"): # if we are the last peer, try and close the iterator
iterator.close() if not peers and hasattr(iterator, "close"):
iterator.close()
class Tee(Generic[T]): class Tee(Generic[T]):