2017-05-11 16:39:54 +02:00
|
|
|
// Copyright 2014 Google Inc. All Rights Reserved.
|
|
|
|
//
|
|
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
|
|
// you may not use this file except in compliance with the License.
|
|
|
|
// You may obtain a copy of the License at
|
|
|
|
//
|
|
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
//
|
|
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
|
|
// See the License for the specific language governing permissions and
|
|
|
|
// limitations under the License.
|
|
|
|
|
|
|
|
package pubsub // import "cloud.google.com/go/pubsub"
|
|
|
|
|
|
|
|
import (
|
|
|
|
"fmt"
|
|
|
|
"os"
|
2017-07-23 09:51:42 +02:00
|
|
|
"runtime"
|
2017-09-30 16:27:27 +02:00
|
|
|
"time"
|
2017-05-11 16:39:54 +02:00
|
|
|
|
2018-03-19 16:51:38 +01:00
|
|
|
"cloud.google.com/go/internal/version"
|
|
|
|
vkit "cloud.google.com/go/pubsub/apiv1"
|
|
|
|
"golang.org/x/net/context"
|
2017-05-11 16:39:54 +02:00
|
|
|
"google.golang.org/api/option"
|
|
|
|
"google.golang.org/grpc"
|
2017-09-30 16:27:27 +02:00
|
|
|
"google.golang.org/grpc/keepalive"
|
2017-05-11 16:39:54 +02:00
|
|
|
)
|
|
|
|
|
|
|
|
const (
|
|
|
|
// ScopePubSub grants permissions to view and manage Pub/Sub
|
|
|
|
// topics and subscriptions.
|
|
|
|
ScopePubSub = "https://www.googleapis.com/auth/pubsub"
|
|
|
|
|
|
|
|
// ScopeCloudPlatform grants permissions to view and manage your data
|
|
|
|
// across Google Cloud Platform services.
|
|
|
|
ScopeCloudPlatform = "https://www.googleapis.com/auth/cloud-platform"
|
|
|
|
)
|
|
|
|
|
2018-05-02 18:09:45 +02:00
|
|
|
const (
|
|
|
|
prodAddr = "https://pubsub.googleapis.com/"
|
|
|
|
minAckDeadline = 10 * time.Second
|
|
|
|
maxAckDeadline = 10 * time.Minute
|
|
|
|
)
|
2017-05-11 16:39:54 +02:00
|
|
|
|
|
|
|
// Client is a Google Pub/Sub client scoped to a single project.
|
|
|
|
//
|
|
|
|
// Clients should be reused rather than being created as needed.
|
|
|
|
// A Client may be shared by multiple goroutines.
|
|
|
|
type Client struct {
|
|
|
|
projectID string
|
2018-03-19 16:51:38 +01:00
|
|
|
pubc *vkit.PublisherClient
|
|
|
|
subc *vkit.SubscriberClient
|
2017-05-11 16:39:54 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
// NewClient creates a new PubSub client.
|
2018-03-19 16:51:38 +01:00
|
|
|
func NewClient(ctx context.Context, projectID string, opts ...option.ClientOption) (c *Client, err error) {
|
2017-05-11 16:39:54 +02:00
|
|
|
var o []option.ClientOption
|
|
|
|
// Environment variables for gcloud emulator:
|
|
|
|
// https://cloud.google.com/sdk/gcloud/reference/beta/emulators/pubsub/
|
|
|
|
if addr := os.Getenv("PUBSUB_EMULATOR_HOST"); addr != "" {
|
|
|
|
conn, err := grpc.Dial(addr, grpc.WithInsecure())
|
|
|
|
if err != nil {
|
|
|
|
return nil, fmt.Errorf("grpc.Dial: %v", err)
|
|
|
|
}
|
|
|
|
o = []option.ClientOption{option.WithGRPCConn(conn)}
|
|
|
|
} else {
|
2017-07-23 09:51:42 +02:00
|
|
|
o = []option.ClientOption{
|
|
|
|
// Create multiple connections to increase throughput.
|
|
|
|
option.WithGRPCConnectionPool(runtime.GOMAXPROCS(0)),
|
2017-09-30 16:27:27 +02:00
|
|
|
option.WithGRPCDialOption(grpc.WithKeepaliveParams(keepalive.ClientParameters{
|
|
|
|
Time: 5 * time.Minute,
|
|
|
|
})),
|
2017-07-23 09:51:42 +02:00
|
|
|
}
|
2018-03-19 16:51:38 +01:00
|
|
|
o = append(o, openCensusOptions()...)
|
2017-05-11 16:39:54 +02:00
|
|
|
}
|
|
|
|
o = append(o, opts...)
|
2018-03-19 16:51:38 +01:00
|
|
|
pubc, err := vkit.NewPublisherClient(ctx, o...)
|
2017-05-11 16:39:54 +02:00
|
|
|
if err != nil {
|
2018-03-19 16:51:38 +01:00
|
|
|
return nil, fmt.Errorf("pubsub: %v", err)
|
2017-05-11 16:39:54 +02:00
|
|
|
}
|
2018-03-19 16:51:38 +01:00
|
|
|
subc, err := vkit.NewSubscriberClient(ctx, option.WithGRPCConn(pubc.Connection()))
|
|
|
|
if err != nil {
|
|
|
|
// Should never happen, since we are passing in the connection.
|
|
|
|
// If it does, we cannot close, because the user may have passed in their
|
|
|
|
// own connection originally.
|
|
|
|
return nil, fmt.Errorf("pubsub: %v", err)
|
2017-05-11 16:39:54 +02:00
|
|
|
}
|
2018-03-19 16:51:38 +01:00
|
|
|
pubc.SetGoogleClientInfo("gccl", version.Repo)
|
|
|
|
subc.SetGoogleClientInfo("gccl", version.Repo)
|
|
|
|
return &Client{
|
|
|
|
projectID: projectID,
|
|
|
|
pubc: pubc,
|
|
|
|
subc: subc,
|
|
|
|
}, nil
|
2017-05-11 16:39:54 +02:00
|
|
|
}
|
|
|
|
|
2018-01-16 14:20:59 +01:00
|
|
|
// Close releases any resources held by the client,
|
|
|
|
// such as memory and goroutines.
|
2017-05-11 16:39:54 +02:00
|
|
|
//
|
2018-01-16 14:20:59 +01:00
|
|
|
// If the client is available for the lifetime of the program, then Close need not be
|
|
|
|
// called at exit.
|
2017-05-11 16:39:54 +02:00
|
|
|
func (c *Client) Close() error {
|
2018-03-19 16:51:38 +01:00
|
|
|
// Return the first error, because the first call closes the connection.
|
|
|
|
err := c.pubc.Close()
|
|
|
|
_ = c.subc.Close()
|
|
|
|
return err
|
2017-05-11 16:39:54 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Client) fullyQualifiedProjectName() string {
|
|
|
|
return fmt.Sprintf("projects/%s", c.projectID)
|
|
|
|
}
|