googleapis / googleapis/google-cloud-python

`StreamingPullManager._shutdown()` fails with `object has no attribute 'nack'` on client shutdown

Aperta
#15,636 3 commenti 1 reazione 1 assegnatario Rivendicata da @abbrowne126 Vedi su GitHub
api: pubsub
Lingua principale
Python
Stelle
5.4k
Fork
1.8k
Merge medio
3g 4h
PR unite (30g)
122

Descrizione

#### Environment details

- OS type and version: Docker image python:3.12-bookworm
- Python version: 3.12.1
- pip version: 23.2.1
- `google-cloud-pubsub` version: 2.31.1
#### Steps to reproduce

1. Subscribe to a PubSub subscription (which returns a future)
2. Call `cancel()` and `result()` on the future to initiate a shutdown sequence

#### Example code
I've extracted below simplified code samples. Those are meant as illustration, they can't be run as is.
```python
from concurrent.futures import ThreadPoolExecutor
from google.cloud import pubsub_v1
from google.cloud.pubsub_v1.subscriber.scheduler import ThreadScheduler

# Message callback
def callback(message):
try:
# ... do something useful ...
message.ack()
except Exception:
message.nack()

# Subscribe to Pubsub subscription
client = pubsub_v1.SubscriberClient()

flow_control = pubsub_v1.types.FlowControl(
max_bytes=10 * 1024 * 1024,
max_messages=100,
max_lease_duration=240 * 60 * 60
)
scheduler = ThreadScheduler(ThreadPoolExecutor(max_workers=10))
future = client.subscribe("subscription_name", callback, flow_control=flow_control, scheduler=scheduler)

# Later, initiate shutdown sequence
future.cancel()
future.result()
```

#### Investigation
There are 2 exceptions from 2 concurrent threads that are logged at the same time.

One is:
```
Exception in thread Thread-RegularStreamShutdown:
Traceback (most recent call last):
File "/usr/local/lib/python3.12/threading.py", line 1073, in _bootstrap_inner
self.run()
File "/usr/local/lib/python3.12/threading.py", line 1010, in run
self._target(*self._args, **self._kwargs)
File "/usr/local/lib/python3.12/site-packages/google/cloud/pubsub_v1/subscriber/_protocol/streaming_pull_manager.py", line 994, in _shutdown
msg.nack()
AttributeError: 'tuple' object has no attribute 'nack'
```

The other is:
```
Exception in thread Thread-OnRpcTerminated:
Traceback (most recent call last):
File "/usr/local/lib/python3.12/threading.py", line 1073, in _bootstrap_inner
self.run()
File "/usr/local/lib/python3.12/threading.py", line 1010, in run
self._target(*self._args, **self._kwargs)
File "/usr/local/lib/python3.12/site-packages/google/cloud/pubsub_v1/subscriber/_protocol/streaming_pull_manager.py", line 967, in _shutdown
assert self._scheduler is not None
AssertionError
```

This looks very much like [another reported bug](https://github.com/googleapis/python-pubsub/issues/997), that's supposed to be fixed as of release 2.23.1, yet I'm running 2.31.1 . Note that [this specific message](https://github.com/googleapis/python-pubsub/issues/997#issuecomment-1898928454) mentioned a similar exception on `msg.nack()`, except it was `NoneType` instead of `tuple`.

What seems to happen is that:
1. first, thread `Thread-RegularStreamShutdown` calls `StreamingPullManager._shutdown()`
- acquires lock `StreamingPullManager._closing`
- asserts that `self._scheduler` is not `None` => OK ✅
- stops `self._scheduler` and resets it to `None`
- tries to `msg.nack()` but fails because `msg` is unexpectedly a `tuple` => raises an exception, releases lock `StreamingPullManager._closing`
2. then, thread `Thread-OnRpcTerminated` also calls `StreamingPullManager._shutdown()` (or, possibly, was already called earlier but waited for the lock to be released)
- acquires lock `StreamingPullManager._closing`
- asserts that `self._scheduler` is not `None` => failure ❌

Thus, the root cause seems to be `msg` being a `tuple`. That variable comes from private fields within the library, hence I suspect an issue in the library code rather than in the user code, but I might be wrong if references to that variable leak to user code.

Guida per i contributori

Apri la guida per i contributori

Valutazione

Questa issue non è ancora stata valutata.

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.