2017-05-11 14:39:54 +00:00
/ *
Copyright 2015 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 bigtable
import (
"fmt"
2018-03-19 15:51:38 +00:00
"math"
2017-05-11 14:39:54 +00:00
"regexp"
"strings"
2018-03-19 15:51:38 +00:00
"time"
2017-05-11 14:39:54 +00:00
2018-03-19 15:51:38 +00:00
"cloud.google.com/go/bigtable/internal/gax"
2017-05-11 14:39:54 +00:00
btopt "cloud.google.com/go/bigtable/internal/option"
"cloud.google.com/go/longrunning"
2017-07-23 07:51:42 +00:00
lroauto "cloud.google.com/go/longrunning/autogen"
2018-03-19 15:51:38 +00:00
"github.com/golang/protobuf/ptypes"
durpb "github.com/golang/protobuf/ptypes/duration"
2017-05-11 14:39:54 +00:00
"golang.org/x/net/context"
2018-03-19 15:51:38 +00:00
"google.golang.org/api/cloudresourcemanager/v1"
"google.golang.org/api/iterator"
2017-05-11 14:39:54 +00:00
"google.golang.org/api/option"
2017-09-30 14:27:27 +00:00
gtransport "google.golang.org/api/transport/grpc"
2017-05-11 14:39:54 +00:00
btapb "google.golang.org/genproto/googleapis/bigtable/admin/v2"
"google.golang.org/grpc"
2018-03-19 15:51:38 +00:00
"google.golang.org/grpc/codes"
2017-05-11 14:39:54 +00:00
"google.golang.org/grpc/metadata"
2017-09-30 14:27:27 +00:00
"google.golang.org/grpc/status"
2017-05-11 14:39:54 +00:00
)
const adminAddr = "bigtableadmin.googleapis.com:443"
// AdminClient is a client type for performing admin operations within a specific instance.
type AdminClient struct {
2018-03-19 15:51:38 +00:00
conn * grpc . ClientConn
tClient btapb . BigtableTableAdminClient
lroClient * lroauto . OperationsClient
2017-05-11 14:39:54 +00:00
project , instance string
// Metadata to be sent with each request.
md metadata . MD
}
// NewAdminClient creates a new AdminClient for a given project and instance.
func NewAdminClient ( ctx context . Context , project , instance string , opts ... option . ClientOption ) ( * AdminClient , error ) {
o , err := btopt . DefaultClientOptions ( adminAddr , AdminScope , clientUserAgent )
if err != nil {
return nil , err
}
2018-03-19 15:51:38 +00:00
// Need to add scopes for long running operations (for create table & snapshots)
o = append ( o , option . WithScopes ( cloudresourcemanager . CloudPlatformScope ) )
2017-05-11 14:39:54 +00:00
o = append ( o , opts ... )
2017-09-30 14:27:27 +00:00
conn , err := gtransport . Dial ( ctx , o ... )
2017-05-11 14:39:54 +00:00
if err != nil {
return nil , fmt . Errorf ( "dialing: %v" , err )
}
2018-03-19 15:51:38 +00:00
lroClient , err := lroauto . NewOperationsClient ( ctx , option . WithGRPCConn ( conn ) )
if err != nil {
// This error "should not happen", since we are just reusing old connection
// and never actually need to dial.
// If this does happen, we could leak conn. However, we cannot close conn:
// If the user invoked the function with option.WithGRPCConn,
// we would close a connection that's still in use.
// TODO(pongad): investigate error conditions.
return nil , err
}
2017-05-11 14:39:54 +00:00
return & AdminClient {
2018-03-19 15:51:38 +00:00
conn : conn ,
tClient : btapb . NewBigtableTableAdminClient ( conn ) ,
lroClient : lroClient ,
project : project ,
instance : instance ,
md : metadata . Pairs ( resourcePrefixHeader , fmt . Sprintf ( "projects/%s/instances/%s" , project , instance ) ) ,
2017-05-11 14:39:54 +00:00
} , nil
}
// Close closes the AdminClient.
func ( ac * AdminClient ) Close ( ) error {
return ac . conn . Close ( )
}
func ( ac * AdminClient ) instancePrefix ( ) string {
return fmt . Sprintf ( "projects/%s/instances/%s" , ac . project , ac . instance )
}
// Tables returns a list of the tables in the instance.
func ( ac * AdminClient ) Tables ( ctx context . Context ) ( [ ] string , error ) {
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , ac . md )
2017-05-11 14:39:54 +00:00
prefix := ac . instancePrefix ( )
req := & btapb . ListTablesRequest {
Parent : prefix ,
}
2018-05-02 16:09:45 +00:00
var res * btapb . ListTablesResponse
err := gax . Invoke ( ctx , func ( ctx context . Context ) error {
var err error
res , err = ac . tClient . ListTables ( ctx , req )
return err
} , retryOptions ... )
2017-05-11 14:39:54 +00:00
if err != nil {
return nil , err
}
2018-05-02 16:09:45 +00:00
2017-05-11 14:39:54 +00:00
names := make ( [ ] string , 0 , len ( res . Tables ) )
for _ , tbl := range res . Tables {
names = append ( names , strings . TrimPrefix ( tbl . Name , prefix + "/tables/" ) )
}
return names , nil
}
2017-09-30 14:27:27 +00:00
// TableConf contains all of the information necessary to create a table with column families.
type TableConf struct {
TableID string
SplitKeys [ ] string
// Families is a map from family name to GCPolicy
Families map [ string ] GCPolicy
}
2017-05-11 14:39:54 +00:00
// CreateTable creates a new table in the instance.
// This method may return before the table's creation is complete.
func ( ac * AdminClient ) CreateTable ( ctx context . Context , table string ) error {
2017-09-30 14:27:27 +00:00
return ac . CreateTableFromConf ( ctx , & TableConf { TableID : table } )
2017-05-11 14:39:54 +00:00
}
// CreatePresplitTable creates a new table in the instance.
// The list of row keys will be used to initially split the table into multiple tablets.
// Given two split keys, "s1" and "s2", three tablets will be created,
// spanning the key ranges: [, s1), [s1, s2), [s2, ).
// This method may return before the table's creation is complete.
2017-09-30 14:27:27 +00:00
func ( ac * AdminClient ) CreatePresplitTable ( ctx context . Context , table string , splitKeys [ ] string ) error {
return ac . CreateTableFromConf ( ctx , & TableConf { TableID : table , SplitKeys : splitKeys } )
}
// CreateTableFromConf creates a new table in the instance from the given configuration.
func ( ac * AdminClient ) CreateTableFromConf ( ctx context . Context , conf * TableConf ) error {
ctx = mergeOutgoingMetadata ( ctx , ac . md )
2017-05-11 14:39:54 +00:00
var req_splits [ ] * btapb . CreateTableRequest_Split
2017-09-30 14:27:27 +00:00
for _ , split := range conf . SplitKeys {
2018-03-19 15:51:38 +00:00
req_splits = append ( req_splits , & btapb . CreateTableRequest_Split { Key : [ ] byte ( split ) } )
2017-05-11 14:39:54 +00:00
}
2017-09-30 14:27:27 +00:00
var tbl btapb . Table
if conf . Families != nil {
tbl . ColumnFamilies = make ( map [ string ] * btapb . ColumnFamily )
for fam , policy := range conf . Families {
2018-03-19 15:51:38 +00:00
tbl . ColumnFamilies [ fam ] = & btapb . ColumnFamily { GcRule : policy . proto ( ) }
2017-09-30 14:27:27 +00:00
}
}
2017-05-11 14:39:54 +00:00
prefix := ac . instancePrefix ( )
req := & btapb . CreateTableRequest {
Parent : prefix ,
2017-09-30 14:27:27 +00:00
TableId : conf . TableID ,
Table : & tbl ,
2017-05-11 14:39:54 +00:00
InitialSplits : req_splits ,
}
_ , err := ac . tClient . CreateTable ( ctx , req )
return err
}
// CreateColumnFamily creates a new column family in a table.
func ( ac * AdminClient ) CreateColumnFamily ( ctx context . Context , table , family string ) error {
// TODO(dsymonds): Permit specifying gcexpr and any other family settings.
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , ac . md )
2017-05-11 14:39:54 +00:00
prefix := ac . instancePrefix ( )
req := & btapb . ModifyColumnFamiliesRequest {
Name : prefix + "/tables/" + table ,
Modifications : [ ] * btapb . ModifyColumnFamiliesRequest_Modification { {
Id : family ,
2018-03-19 15:51:38 +00:00
Mod : & btapb . ModifyColumnFamiliesRequest_Modification_Create { Create : & btapb . ColumnFamily { } } ,
2017-05-11 14:39:54 +00:00
} } ,
}
_ , err := ac . tClient . ModifyColumnFamilies ( ctx , req )
return err
}
// DeleteTable deletes a table and all of its data.
func ( ac * AdminClient ) DeleteTable ( ctx context . Context , table string ) error {
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , ac . md )
2017-05-11 14:39:54 +00:00
prefix := ac . instancePrefix ( )
req := & btapb . DeleteTableRequest {
Name : prefix + "/tables/" + table ,
}
_ , err := ac . tClient . DeleteTable ( ctx , req )
return err
}
// DeleteColumnFamily deletes a column family in a table and all of its data.
func ( ac * AdminClient ) DeleteColumnFamily ( ctx context . Context , table , family string ) error {
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , ac . md )
2017-05-11 14:39:54 +00:00
prefix := ac . instancePrefix ( )
req := & btapb . ModifyColumnFamiliesRequest {
Name : prefix + "/tables/" + table ,
Modifications : [ ] * btapb . ModifyColumnFamiliesRequest_Modification { {
Id : family ,
2018-03-19 15:51:38 +00:00
Mod : & btapb . ModifyColumnFamiliesRequest_Modification_Drop { Drop : true } ,
2017-05-11 14:39:54 +00:00
} } ,
}
_ , err := ac . tClient . ModifyColumnFamilies ( ctx , req )
return err
}
// TableInfo represents information about a table.
type TableInfo struct {
2017-07-23 07:51:42 +00:00
// DEPRECATED - This field is deprecated. Please use FamilyInfos instead.
2017-09-30 14:27:27 +00:00
Families [ ] string
2017-07-23 07:51:42 +00:00
FamilyInfos [ ] FamilyInfo
}
// FamilyInfo represents information about a column family.
type FamilyInfo struct {
2017-09-30 14:27:27 +00:00
Name string
2017-07-23 07:51:42 +00:00
GCPolicy string
2017-05-11 14:39:54 +00:00
}
// TableInfo retrieves information about a table.
func ( ac * AdminClient ) TableInfo ( ctx context . Context , table string ) ( * TableInfo , error ) {
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , ac . md )
2017-05-11 14:39:54 +00:00
prefix := ac . instancePrefix ( )
req := & btapb . GetTableRequest {
Name : prefix + "/tables/" + table ,
}
2018-05-02 16:09:45 +00:00
var res * btapb . Table
err := gax . Invoke ( ctx , func ( ctx context . Context ) error {
var err error
res , err = ac . tClient . GetTable ( ctx , req )
return err
} , retryOptions ... )
2017-05-11 14:39:54 +00:00
if err != nil {
return nil , err
}
2018-05-02 16:09:45 +00:00
2017-05-11 14:39:54 +00:00
ti := & TableInfo { }
2017-07-23 07:51:42 +00:00
for name , fam := range res . ColumnFamilies {
ti . Families = append ( ti . Families , name )
ti . FamilyInfos = append ( ti . FamilyInfos , FamilyInfo { Name : name , GCPolicy : GCRuleToString ( fam . GcRule ) } )
2017-05-11 14:39:54 +00:00
}
return ti , nil
}
// SetGCPolicy specifies which cells in a column family should be garbage collected.
// GC executes opportunistically in the background; table reads may return data
// matching the GC policy.
func ( ac * AdminClient ) SetGCPolicy ( ctx context . Context , table , family string , policy GCPolicy ) error {
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , ac . md )
2017-05-11 14:39:54 +00:00
prefix := ac . instancePrefix ( )
req := & btapb . ModifyColumnFamiliesRequest {
Name : prefix + "/tables/" + table ,
Modifications : [ ] * btapb . ModifyColumnFamiliesRequest_Modification { {
Id : family ,
2018-03-19 15:51:38 +00:00
Mod : & btapb . ModifyColumnFamiliesRequest_Modification_Update { Update : & btapb . ColumnFamily { GcRule : policy . proto ( ) } } ,
2017-05-11 14:39:54 +00:00
} } ,
}
_ , err := ac . tClient . ModifyColumnFamilies ( ctx , req )
return err
}
2017-07-23 07:51:42 +00:00
// DropRowRange permanently deletes a row range from the specified table.
func ( ac * AdminClient ) DropRowRange ( ctx context . Context , table , rowKeyPrefix string ) error {
ctx = mergeOutgoingMetadata ( ctx , ac . md )
prefix := ac . instancePrefix ( )
req := & btapb . DropRowRangeRequest {
Name : prefix + "/tables/" + table ,
2018-03-19 15:51:38 +00:00
Target : & btapb . DropRowRangeRequest_RowKeyPrefix { RowKeyPrefix : [ ] byte ( rowKeyPrefix ) } ,
2017-07-23 07:51:42 +00:00
}
_ , err := ac . tClient . DropRowRange ( ctx , req )
return err
}
2018-03-19 15:51:38 +00:00
// CreateTableFromSnapshot creates a table from snapshot.
// The table will be created in the same cluster as the snapshot.
//
// This is a private alpha release of Cloud Bigtable snapshots. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( ac * AdminClient ) CreateTableFromSnapshot ( ctx context . Context , table , cluster , snapshot string ) error {
ctx = mergeOutgoingMetadata ( ctx , ac . md )
prefix := ac . instancePrefix ( )
snapshotPath := prefix + "/clusters/" + cluster + "/snapshots/" + snapshot
req := & btapb . CreateTableFromSnapshotRequest {
Parent : prefix ,
TableId : table ,
SourceSnapshot : snapshotPath ,
}
op , err := ac . tClient . CreateTableFromSnapshot ( ctx , req )
if err != nil {
return err
}
resp := btapb . Table { }
return longrunning . InternalNewOperation ( ac . lroClient , op ) . Wait ( ctx , & resp )
}
const DefaultSnapshotDuration time . Duration = 0
// Creates a new snapshot in the specified cluster from the specified source table.
// Setting the ttl to `DefaultSnapshotDuration` will use the server side default for the duration.
//
// This is a private alpha release of Cloud Bigtable snapshots. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( ac * AdminClient ) SnapshotTable ( ctx context . Context , table , cluster , snapshot string , ttl time . Duration ) error {
ctx = mergeOutgoingMetadata ( ctx , ac . md )
prefix := ac . instancePrefix ( )
var ttlProto * durpb . Duration
if ttl > 0 {
ttlProto = ptypes . DurationProto ( ttl )
}
req := & btapb . SnapshotTableRequest {
Name : prefix + "/tables/" + table ,
Cluster : prefix + "/clusters/" + cluster ,
SnapshotId : snapshot ,
Ttl : ttlProto ,
}
op , err := ac . tClient . SnapshotTable ( ctx , req )
if err != nil {
return err
}
resp := btapb . Snapshot { }
return longrunning . InternalNewOperation ( ac . lroClient , op ) . Wait ( ctx , & resp )
}
// Returns a SnapshotIterator for iterating over the snapshots in a cluster.
// To list snapshots across all of the clusters in the instance specify "-" as the cluster.
//
// This is a private alpha release of Cloud Bigtable snapshots. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( ac * AdminClient ) ListSnapshots ( ctx context . Context , cluster string ) * SnapshotIterator {
ctx = mergeOutgoingMetadata ( ctx , ac . md )
prefix := ac . instancePrefix ( )
clusterPath := prefix + "/clusters/" + cluster
it := & SnapshotIterator { }
req := & btapb . ListSnapshotsRequest {
Parent : clusterPath ,
}
fetch := func ( pageSize int , pageToken string ) ( string , error ) {
req . PageToken = pageToken
if pageSize > math . MaxInt32 {
req . PageSize = math . MaxInt32
} else {
req . PageSize = int32 ( pageSize )
}
resp , err := ac . tClient . ListSnapshots ( ctx , req )
if err != nil {
return "" , err
}
for _ , s := range resp . Snapshots {
snapshotInfo , err := newSnapshotInfo ( s )
if err != nil {
return "" , fmt . Errorf ( "Failed to parse snapshot proto %v" , err )
}
it . items = append ( it . items , snapshotInfo )
}
return resp . NextPageToken , nil
}
bufLen := func ( ) int { return len ( it . items ) }
takeBuf := func ( ) interface { } { b := it . items ; it . items = nil ; return b }
it . pageInfo , it . nextFunc = iterator . NewPageInfo ( fetch , bufLen , takeBuf )
return it
}
func newSnapshotInfo ( snapshot * btapb . Snapshot ) ( * SnapshotInfo , error ) {
nameParts := strings . Split ( snapshot . Name , "/" )
name := nameParts [ len ( nameParts ) - 1 ]
tablePathParts := strings . Split ( snapshot . SourceTable . Name , "/" )
tableId := tablePathParts [ len ( tablePathParts ) - 1 ]
createTime , err := ptypes . Timestamp ( snapshot . CreateTime )
if err != nil {
return nil , fmt . Errorf ( "Invalid createTime: %v" , err )
}
deleteTime , err := ptypes . Timestamp ( snapshot . DeleteTime )
if err != nil {
return nil , fmt . Errorf ( "Invalid deleteTime: %v" , err )
}
return & SnapshotInfo {
Name : name ,
SourceTable : tableId ,
DataSize : snapshot . DataSizeBytes ,
CreateTime : createTime ,
DeleteTime : deleteTime ,
} , nil
}
// An EntryIterator iterates over log entries.
//
// This is a private alpha release of Cloud Bigtable snapshots. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
type SnapshotIterator struct {
items [ ] * SnapshotInfo
pageInfo * iterator . PageInfo
nextFunc func ( ) error
}
// PageInfo supports pagination. See https://godoc.org/google.golang.org/api/iterator package for details.
func ( it * SnapshotIterator ) PageInfo ( ) * iterator . PageInfo {
return it . pageInfo
}
// Next returns the next result. Its second return value is iterator.Done
// (https://godoc.org/google.golang.org/api/iterator) if there are no more
// results. Once Next returns Done, all subsequent calls will return Done.
func ( it * SnapshotIterator ) Next ( ) ( * SnapshotInfo , error ) {
if err := it . nextFunc ( ) ; err != nil {
return nil , err
}
item := it . items [ 0 ]
it . items = it . items [ 1 : ]
return item , nil
}
type SnapshotInfo struct {
Name string
SourceTable string
DataSize int64
CreateTime time . Time
DeleteTime time . Time
}
// Get snapshot metadata.
//
// This is a private alpha release of Cloud Bigtable snapshots. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( ac * AdminClient ) SnapshotInfo ( ctx context . Context , cluster , snapshot string ) ( * SnapshotInfo , error ) {
ctx = mergeOutgoingMetadata ( ctx , ac . md )
prefix := ac . instancePrefix ( )
clusterPath := prefix + "/clusters/" + cluster
snapshotPath := clusterPath + "/snapshots/" + snapshot
req := & btapb . GetSnapshotRequest {
Name : snapshotPath ,
}
resp , err := ac . tClient . GetSnapshot ( ctx , req )
if err != nil {
return nil , err
}
return newSnapshotInfo ( resp )
}
// Delete a snapshot in a cluster.
//
// This is a private alpha release of Cloud Bigtable snapshots. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( ac * AdminClient ) DeleteSnapshot ( ctx context . Context , cluster , snapshot string ) error {
ctx = mergeOutgoingMetadata ( ctx , ac . md )
prefix := ac . instancePrefix ( )
clusterPath := prefix + "/clusters/" + cluster
snapshotPath := clusterPath + "/snapshots/" + snapshot
req := & btapb . DeleteSnapshotRequest {
Name : snapshotPath ,
}
_ , err := ac . tClient . DeleteSnapshot ( ctx , req )
return err
}
// getConsistencyToken gets the consistency token for a table.
func ( ac * AdminClient ) getConsistencyToken ( ctx context . Context , tableName string ) ( string , error ) {
req := & btapb . GenerateConsistencyTokenRequest {
Name : tableName ,
}
resp , err := ac . tClient . GenerateConsistencyToken ( ctx , req )
if err != nil {
return "" , err
}
return resp . GetConsistencyToken ( ) , nil
}
// isConsistent checks if a token is consistent for a table.
func ( ac * AdminClient ) isConsistent ( ctx context . Context , tableName , token string ) ( bool , error ) {
req := & btapb . CheckConsistencyRequest {
Name : tableName ,
ConsistencyToken : token ,
}
var resp * btapb . CheckConsistencyResponse
// Retry calls on retryable errors to avoid losing the token gathered before.
err := gax . Invoke ( ctx , func ( ctx context . Context ) error {
var err error
resp , err = ac . tClient . CheckConsistency ( ctx , req )
return err
} , retryOptions ... )
if err != nil {
return false , err
}
return resp . GetConsistent ( ) , nil
}
// WaitForReplication waits until all the writes committed before the call started have been propagated to all the clusters in the instance via replication.
//
// This is a private alpha release of Cloud Bigtable replication. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( ac * AdminClient ) WaitForReplication ( ctx context . Context , table string ) error {
// Get the token.
prefix := ac . instancePrefix ( )
tableName := prefix + "/tables/" + table
token , err := ac . getConsistencyToken ( ctx , tableName )
if err != nil {
return err
}
// Periodically check if the token is consistent.
timer := time . NewTicker ( time . Second * 10 )
defer timer . Stop ( )
for {
consistent , err := ac . isConsistent ( ctx , tableName , token )
if err != nil {
return err
}
if consistent {
return nil
}
// Sleep for a bit or until the ctx is cancelled.
select {
case <- ctx . Done ( ) :
return ctx . Err ( )
case <- timer . C :
}
}
}
2017-05-11 14:39:54 +00:00
const instanceAdminAddr = "bigtableadmin.googleapis.com:443"
// InstanceAdminClient is a client type for performing admin operations on instances.
// These operations can be substantially more dangerous than those provided by AdminClient.
type InstanceAdminClient struct {
2017-07-23 07:51:42 +00:00
conn * grpc . ClientConn
iClient btapb . BigtableInstanceAdminClient
lroClient * lroauto . OperationsClient
2017-05-11 14:39:54 +00:00
project string
// Metadata to be sent with each request.
md metadata . MD
}
// NewInstanceAdminClient creates a new InstanceAdminClient for a given project.
func NewInstanceAdminClient ( ctx context . Context , project string , opts ... option . ClientOption ) ( * InstanceAdminClient , error ) {
o , err := btopt . DefaultClientOptions ( instanceAdminAddr , InstanceAdminScope , clientUserAgent )
if err != nil {
return nil , err
}
o = append ( o , opts ... )
2017-09-30 14:27:27 +00:00
conn , err := gtransport . Dial ( ctx , o ... )
2017-05-11 14:39:54 +00:00
if err != nil {
return nil , fmt . Errorf ( "dialing: %v" , err )
}
2017-07-23 07:51:42 +00:00
lroClient , err := lroauto . NewOperationsClient ( ctx , option . WithGRPCConn ( conn ) )
if err != nil {
// This error "should not happen", since we are just reusing old connection
// and never actually need to dial.
// If this does happen, we could leak conn. However, we cannot close conn:
// If the user invoked the function with option.WithGRPCConn,
// we would close a connection that's still in use.
// TODO(pongad): investigate error conditions.
return nil , err
}
2017-05-11 14:39:54 +00:00
return & InstanceAdminClient {
2017-07-23 07:51:42 +00:00
conn : conn ,
iClient : btapb . NewBigtableInstanceAdminClient ( conn ) ,
lroClient : lroClient ,
2017-05-11 14:39:54 +00:00
project : project ,
md : metadata . Pairs ( resourcePrefixHeader , "projects/" + project ) ,
} , nil
}
// Close closes the InstanceAdminClient.
func ( iac * InstanceAdminClient ) Close ( ) error {
return iac . conn . Close ( )
}
// StorageType is the type of storage used for all tables in an instance
type StorageType int
const (
SSD StorageType = iota
HDD
)
func ( st StorageType ) proto ( ) btapb . StorageType {
if st == HDD {
return btapb . StorageType_HDD
}
return btapb . StorageType_SSD
}
2017-09-30 14:27:27 +00:00
// InstanceType is the type of the instance
type InstanceType int32
const (
PRODUCTION InstanceType = InstanceType ( btapb . Instance_PRODUCTION )
DEVELOPMENT = InstanceType ( btapb . Instance_DEVELOPMENT )
)
2017-05-11 14:39:54 +00:00
// InstanceInfo represents information about an instance
type InstanceInfo struct {
Name string // name of the instance
DisplayName string // display name for UIs
}
// InstanceConf contains the information necessary to create an Instance
type InstanceConf struct {
InstanceId , DisplayName , ClusterId , Zone string
2017-09-30 14:27:27 +00:00
// NumNodes must not be specified for DEVELOPMENT instance types
NumNodes int32
StorageType StorageType
InstanceType InstanceType
2017-05-11 14:39:54 +00:00
}
2018-03-19 15:51:38 +00:00
// InstanceWithClustersConfig contains the information necessary to create an Instance
type InstanceWithClustersConfig struct {
InstanceID , DisplayName string
Clusters [ ] ClusterConfig
InstanceType InstanceType
}
2017-05-11 14:39:54 +00:00
var instanceNameRegexp = regexp . MustCompile ( ` ^projects/([^/]+)/instances/([a-z][-a-z0-9]*)$ ` )
// CreateInstance creates a new instance in the project.
// This method will return when the instance has been created or when an error occurs.
func ( iac * InstanceAdminClient ) CreateInstance ( ctx context . Context , conf * InstanceConf ) error {
2018-03-19 15:51:38 +00:00
newConfig := InstanceWithClustersConfig {
InstanceID : conf . InstanceId ,
DisplayName : conf . DisplayName ,
InstanceType : conf . InstanceType ,
Clusters : [ ] ClusterConfig {
{
InstanceID : conf . InstanceId ,
ClusterID : conf . ClusterId ,
Zone : conf . Zone ,
NumNodes : conf . NumNodes ,
StorageType : conf . StorageType ,
} ,
} ,
}
return iac . CreateInstanceWithClusters ( ctx , & newConfig )
}
// CreateInstance creates a new instance with configured clusters in the project.
// This method will return when the instance has been created or when an error occurs.
//
// Instances with multiple clusters are part of a private alpha release of Cloud Bigtable replication.
// This feature is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( iac * InstanceAdminClient ) CreateInstanceWithClusters ( ctx context . Context , conf * InstanceWithClustersConfig ) error {
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , iac . md )
2018-03-19 15:51:38 +00:00
clusters := make ( map [ string ] * btapb . Cluster )
for _ , cluster := range conf . Clusters {
clusters [ cluster . ClusterID ] = cluster . proto ( iac . project )
}
2017-05-11 14:39:54 +00:00
req := & btapb . CreateInstanceRequest {
Parent : "projects/" + iac . project ,
2018-03-19 15:51:38 +00:00
InstanceId : conf . InstanceID ,
2017-09-30 14:27:27 +00:00
Instance : & btapb . Instance { DisplayName : conf . DisplayName , Type : btapb . Instance_Type ( conf . InstanceType ) } ,
2018-03-19 15:51:38 +00:00
Clusters : clusters ,
2017-05-11 14:39:54 +00:00
}
lro , err := iac . iClient . CreateInstance ( ctx , req )
if err != nil {
return err
}
resp := btapb . Instance { }
2017-07-23 07:51:42 +00:00
return longrunning . InternalNewOperation ( iac . lroClient , lro ) . Wait ( ctx , & resp )
2017-05-11 14:39:54 +00:00
}
// DeleteInstance deletes an instance from the project.
func ( iac * InstanceAdminClient ) DeleteInstance ( ctx context . Context , instanceId string ) error {
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , iac . md )
2018-03-19 15:51:38 +00:00
req := & btapb . DeleteInstanceRequest { Name : "projects/" + iac . project + "/instances/" + instanceId }
2017-05-11 14:39:54 +00:00
_ , err := iac . iClient . DeleteInstance ( ctx , req )
return err
}
// Instances returns a list of instances in the project.
func ( iac * InstanceAdminClient ) Instances ( ctx context . Context ) ( [ ] * InstanceInfo , error ) {
2017-07-23 07:51:42 +00:00
ctx = mergeOutgoingMetadata ( ctx , iac . md )
2017-05-11 14:39:54 +00:00
req := & btapb . ListInstancesRequest {
Parent : "projects/" + iac . project ,
}
res , err := iac . iClient . ListInstances ( ctx , req )
if err != nil {
return nil , err
}
2017-09-30 14:27:27 +00:00
if len ( res . FailedLocations ) > 0 {
// We don't have a good way to return a partial result in the face of some zones being unavailable.
// Fail the entire request.
return nil , status . Errorf ( codes . Unavailable , "Failed locations: %v" , res . FailedLocations )
}
2017-05-11 14:39:54 +00:00
var is [ ] * InstanceInfo
for _ , i := range res . Instances {
m := instanceNameRegexp . FindStringSubmatch ( i . Name )
if m == nil {
return nil , fmt . Errorf ( "malformed instance name %q" , i . Name )
}
is = append ( is , & InstanceInfo {
Name : m [ 2 ] ,
DisplayName : i . DisplayName ,
} )
}
return is , nil
}
2017-07-23 07:51:42 +00:00
// InstanceInfo returns information about an instance.
func ( iac * InstanceAdminClient ) InstanceInfo ( ctx context . Context , instanceId string ) ( * InstanceInfo , error ) {
ctx = mergeOutgoingMetadata ( ctx , iac . md )
req := & btapb . GetInstanceRequest {
Name : "projects/" + iac . project + "/instances/" + instanceId ,
}
res , err := iac . iClient . GetInstance ( ctx , req )
if err != nil {
return nil , err
}
m := instanceNameRegexp . FindStringSubmatch ( res . Name )
if m == nil {
return nil , fmt . Errorf ( "malformed instance name %q" , res . Name )
}
return & InstanceInfo {
Name : m [ 2 ] ,
DisplayName : res . DisplayName ,
} , nil
}
2018-03-19 15:51:38 +00:00
// ClusterConfig contains the information necessary to create a cluster
type ClusterConfig struct {
InstanceID , ClusterID , Zone string
NumNodes int32
StorageType StorageType
}
func ( cc * ClusterConfig ) proto ( project string ) * btapb . Cluster {
return & btapb . Cluster {
ServeNodes : cc . NumNodes ,
DefaultStorageType : cc . StorageType . proto ( ) ,
Location : "projects/" + project + "/locations/" + cc . Zone ,
}
}
// ClusterInfo represents information about a cluster.
type ClusterInfo struct {
Name string // name of the cluster
Zone string // GCP zone of the cluster (e.g. "us-central1-a")
ServeNodes int // number of allocated serve nodes
State string // state of the cluster
}
// CreateCluster creates a new cluster in an instance.
// This method will return when the cluster has been created or when an error occurs.
//
// This is a private alpha release of Cloud Bigtable replication. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( iac * InstanceAdminClient ) CreateCluster ( ctx context . Context , conf * ClusterConfig ) error {
ctx = mergeOutgoingMetadata ( ctx , iac . md )
req := & btapb . CreateClusterRequest {
Parent : "projects/" + iac . project + "/instances/" + conf . InstanceID ,
ClusterId : conf . ClusterID ,
Cluster : conf . proto ( iac . project ) ,
}
lro , err := iac . iClient . CreateCluster ( ctx , req )
if err != nil {
return err
}
resp := btapb . Cluster { }
return longrunning . InternalNewOperation ( iac . lroClient , lro ) . Wait ( ctx , & resp )
}
// DeleteCluster deletes a cluster from an instance.
//
// This is a private alpha release of Cloud Bigtable replication. This feature
// is not currently available to most Cloud Bigtable customers. This feature
// might be changed in backward-incompatible ways and is not recommended for
// production use. It is not subject to any SLA or deprecation policy.
func ( iac * InstanceAdminClient ) DeleteCluster ( ctx context . Context , instanceId , clusterId string ) error {
ctx = mergeOutgoingMetadata ( ctx , iac . md )
req := & btapb . DeleteClusterRequest { Name : "projects/" + iac . project + "/instances/" + instanceId + "/clusters/" + clusterId }
_ , err := iac . iClient . DeleteCluster ( ctx , req )
return err
}
// UpdateCluster updates attributes of a cluster
func ( iac * InstanceAdminClient ) UpdateCluster ( ctx context . Context , instanceId , clusterId string , serveNodes int32 ) error {
ctx = mergeOutgoingMetadata ( ctx , iac . md )
cluster := & btapb . Cluster {
Name : "projects/" + iac . project + "/instances/" + instanceId + "/clusters/" + clusterId ,
ServeNodes : serveNodes }
lro , err := iac . iClient . UpdateCluster ( ctx , cluster )
if err != nil {
return err
}
return longrunning . InternalNewOperation ( iac . lroClient , lro ) . Wait ( ctx , nil )
}
// Clusters lists the clusters in an instance.
func ( iac * InstanceAdminClient ) Clusters ( ctx context . Context , instanceId string ) ( [ ] * ClusterInfo , error ) {
ctx = mergeOutgoingMetadata ( ctx , iac . md )
req := & btapb . ListClustersRequest { Parent : "projects/" + iac . project + "/instances/" + instanceId }
res , err := iac . iClient . ListClusters ( ctx , req )
if err != nil {
return nil , err
}
// TODO(garyelliott): Deal with failed_locations.
var cis [ ] * ClusterInfo
for _ , c := range res . Clusters {
nameParts := strings . Split ( c . Name , "/" )
locParts := strings . Split ( c . Location , "/" )
cis = append ( cis , & ClusterInfo {
Name : nameParts [ len ( nameParts ) - 1 ] ,
Zone : locParts [ len ( locParts ) - 1 ] ,
ServeNodes : int ( c . ServeNodes ) ,
State : c . State . String ( ) ,
} )
}
return cis , nil
}
// GetCluster fetches a cluster in an instance
func ( iac * InstanceAdminClient ) GetCluster ( ctx context . Context , instanceID , clusterID string ) ( * ClusterInfo , error ) {
ctx = mergeOutgoingMetadata ( ctx , iac . md )
2018-06-17 16:59:12 +00:00
req := & btapb . GetClusterRequest { Name : "projects/" + iac . project + "/instances/" + instanceID + "/clusters/" + clusterID }
2018-03-19 15:51:38 +00:00
c , err := iac . iClient . GetCluster ( ctx , req )
if err != nil {
return nil , err
}
nameParts := strings . Split ( c . Name , "/" )
locParts := strings . Split ( c . Location , "/" )
cis := & ClusterInfo {
Name : nameParts [ len ( nameParts ) - 1 ] ,
Zone : locParts [ len ( locParts ) - 1 ] ,
ServeNodes : int ( c . ServeNodes ) ,
State : c . State . String ( ) ,
}
return cis , nil
}