Skip to content

[AI Improvement] [Task] Clean up PubSub streaming subscriptions and gRPC client in PubsubEmulator.stop - #11142

Open
joehan wants to merge 8 commits into
mainfrom
ai-improve-565073134-task-clean-up-pubsub-streaming-subs
Open

joehan wants to merge 8 commits into
mainfrom
ai-improve-565073134-task-clean-up-pubsub-streaming-subs

Conversation

@joehan

@joehan joehan commented Sep 22, 2026

Copy link
Copy Markdown
Member

Description

Resolves Buganizer b/565073134

In src/emulator/pubsubEmulator.ts, this.subscriptionForTopic maintains active Google Cloud Pub/Sub Subscription instances connected to the emulator. When PubsubEmulator.stop() was invoked, it only stopped the underlying downloadable emulator process, leaving active StreamingPull gRPC channels open. As a result, the orphaned gRPC client repeatedly attempted to reconnect to the closed port, logging UNAVAILABLE errors, leaking socket handles and timers, and keeping the Node.js event loop active.

This change:

  1. Iterates over all active subscriptions in this.subscriptionForTopic.values() and closes each cleanly before shutting down the emulator process.
  2. Catches and logs individual subscription close errors gracefully to ensure shutdown proceeds even if an individual connection has already severed.
  3. Closes the underlying PubSub gRPC client via await this._pubsub.close() and resets this._pubsub = undefined.
  4. Clears this.subscriptionForTopic and this.triggersForTopic.
  5. Guarantees that downloadableEmulators.stop(Emulators.PUBSUB) is always executed in a finally block.
  6. Adds unit tests in src/emulator/pubsubEmulator.spec.ts validating complete teardown and error resilience.

Scenarios Tested

  • Clean shutdown with no subscriptions or client initialized.
  • Clean shutdown with active subscriptions and topic triggers, verifying all subscriptions are closed and internal maps cleared.
  • Clean shutdown with initialized _pubsub client, verifying gRPC client close and property reset.
  • Graceful error handling when individual subscriptions reject during close, verifying all remaining subscriptions and the binary stop proceed.
  • Verified with npm run build and npx mocha 'src/emulator/pubsubEmulator.spec.ts'.

@joehan joehan self-assigned this Sep 22, 2026

@gemini-code-assist gemini-code-assist Bot left a comment

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.

Code Review

This pull request updates the PubsubEmulator's stop method to gracefully close active subscriptions and the Pub/Sub client before stopping the emulator, and adds corresponding unit tests. The feedback highlights multiple instances where JSON.stringify(err) is used on Error objects, which results in empty objects ("{}") in the debug logs; it is recommended to safely extract the error message instead.

Comment thread src/emulator/pubsubEmulator.ts
Comment thread src/emulator/pubsubEmulator.ts
Comment thread src/emulator/pubsubEmulator.ts
Comment thread src/emulator/pubsubEmulator.ts Outdated
joehan and others added 2 commits October 9, 2026 15:13
Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
@joehan
joehan marked this pull request as ready for review October 9, 2026 22:13
Comment thread src/emulator/pubsubEmulator.ts Outdated
this.logger.logLabeled("DEBUG", "pubsub", "Pubsub kill output: " + JSON.stringify(buffer));
const closePromises = Array.from(this.subscriptionForTopic.values()).map(async (sub) => {
try {
await sub.close();

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.

sub.close() has no time limit, so shutdown could hang. In @google-cloud/pubsub v5 (^5.2.0 here), Subscriber.close() nacks in-flight messages and then waits for the ack/nack flush. The default wait is maxExtensionTime, which is 60 minutes by default. In the normal case the emulator is still up when this runs, so it should be quick. But if the emulator process has already died, Ctrl-C on emulators:start or emulators:exec could hang there. Either:

set closeOptions: { timeout: ... } (e.g., a few seconds) when the subscription is created/obtained in maybeCreateTopicAndSub (the field is a subscriber option — check that createSubscription passes it through), or

wrap the close calls in Promise.race with a short timer.

A test that hangs one subscription's close would cover this.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Wrapped all subscription close() calls (and this._pubsub.close()) in a withTimeout helper bounded by SUBSCRIPTION_CLOSE_TIMEOUT_MS (2s) via Promise.race with proper timer cleanup. Added unit tests using fake timers (sandbox.useFakeTimers()) to verify that hanging subscription or client close calls time out and do not block shutdown.

});
await Promise.all(closePromises);
this.subscriptionForTopic.clear();
this.triggersForTopic.clear();

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.

Because await this.maybeCreateTopicAndSub(topicName) is called unconditionally before checking this.subscriptionForTopic.has(topicName) or triggers.some(...), every call to addTrigger() for an existing topicName (e.g., when multiple Cloud Functions trigger on the same Pub/Sub topic, or when triggers are re-registered on reload) creates a new Subscription instance with a live sub.on("message", ...) StreamingPull listener and overwrites the previous Subscription in this.subscriptionForTopic (or drops it on early return).
Those overwritten Subscription instances are lost from this.subscriptionForTopic and therefore won't be closed by stop().

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Updated addTrigger to check if (!this.subscriptionForTopic.has(topicName)) before creating topic/sub. When multiple Cloud Functions trigger on the same topic (or triggers are registered across reloads), the existing subscription is preserved and reused. Added tests in pubsubEmulator.spec.ts verifying that maybeCreateTopicAndSub is called only once for multiple triggers on the same topic and that duplicate triggers are ignored.

Comment thread src/emulator/pubsubEmulator.spec.ts Outdated
const emulator = new PubsubEmulator({ projectId: "test-project" });

const closeStub = sandbox.stub().resolves();
const emulatorWithClient = emulator as unknown as {

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.

Can we remove all references of "as unknown as" where practical?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Defined lightweight structural interfaces PubsubSubscription and PubsubClient in src/emulator/pubsubEmulator.ts so that test doubles satisfy types directly without requiring type assertions. Removed all references of as unknown as (and as any) from pubsubEmulator.spec.ts.

This branch has not been deployed

No deployments
Sign up for free to 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