worker

package
v1.9.11-63c34d2a5f789d... Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jan 30, 2020 License: Apache-2.0 Imports: 73 Imported by: 10

Documentation

Index

Constants

View Source
const (
	// PrometheusPort is the port the aggregated metrics are served on for scraping
	PrometheusPort = 9090
)
View Source
const (
	// WorkerEtcdPrefix is the prefix in etcd that we use to store worker information.
	WorkerEtcdPrefix = "workers"
)

Variables

View Source
var (
	ErrInvalidLengthWorkerService        = fmt.Errorf("proto: negative length found during unmarshaling")
	ErrIntOverflowWorkerService          = fmt.Errorf("proto: integer overflow")
	ErrUnexpectedEndOfGroupWorkerService = fmt.Errorf("proto: unexpected end of group")
)
View Source
var State_name = map[int32]string{
	0: "RUNNING",
	1: "COMPLETE",
	3: "FAILED",
}
View Source
var State_value = map[string]int32{
	"RUNNING":  0,
	"COMPLETE": 1,
	"FAILED":   3,
}

Functions

func Cancel added in v1.7.5

func Cancel(ctx context.Context, pipelineRcName string, etcdClient *etcd.Client,
	etcdPrefix string, workerGrpcPort uint16, jobID string, dataFilter []string) error

Cancel cancels a set of datums running on workers. pipelineRcName is the name of the pipeline's RC and can be gotten with ppsutil.PipelineRcName.

func Conns added in v1.7.5

func Conns(ctx context.Context, pipelineRcName string, etcdClient *etcd.Client, etcdPrefix string, workerGrpcPort uint16) ([]*grpc.ClientConn, error)

Conns returns a slice of connections to worker servers. pipelineRcName is the name of the pipeline's RC and can be gotten with ppsutil.PipelineRcName. You can also pass "" for pipelineRcName to get all clients for all workers.

func HashDatum

func HashDatum(pipelineName string, pipelineSalt string, data []*Input) string

HashDatum computes and returns the hash of datum + pipeline, with a pipeline-specific prefix.

func HashDatum15

func HashDatum15(pipelineInfo *pps.PipelineInfo, data []*Input) (string, error)

HashDatum15 computes and returns the hash of datum + pipeline for version <= 1.5.0, with a pipeline-specific prefix.

func MatchDatum

func MatchDatum(filter []string, data []*pps.InputFile) bool

MatchDatum checks if a datum matches a filter. To match each string in filter must correspond match at least 1 datum's Path or Hash. Order of filter and data is irrelevant.

func RegisterWorkerServer

func RegisterWorkerServer(s *grpc.Server, srv WorkerServer)

func Status added in v1.7.5

func Status(ctx context.Context, pipelineRcName string, etcdClient *etcd.Client, etcdPrefix string, workerGrpcPort uint16) ([]*pps.WorkerStatus, error)

Status returns the statuses of workers referenced by pipelineRcName. pipelineRcName is the name of the pipeline's RC and can be gotten with ppsutil.PipelineRcName. You can also pass "" for pipelineRcName to get all clients for all workers.

Types

type APIServer

type APIServer struct {
	// contains filtered or unexported fields
}

APIServer implements the worker API

func NewAPIServer

func NewAPIServer(pachClient *client.APIClient, etcdClient *etcd.Client, etcdPrefix string, pipelineInfo *pps.PipelineInfo, workerName string, namespace string, hashtreeStorage string) (*APIServer, error)

NewAPIServer creates an APIServer for a given pipeline

func (*APIServer) Cancel

func (a *APIServer) Cancel(ctx context.Context, request *CancelRequest) (*CancelResponse, error)

Cancel cancels the currently running datum

func (*APIServer) DatumID added in v1.6.0

func (a *APIServer) DatumID(data []*Input) string

DatumID computes the id for a datum, this value is used in ListDatum and InspectDatum.

func (*APIServer) GetChunk added in v1.8.5

func (a *APIServer) GetChunk(request *GetChunkRequest, server Worker_GetChunkServer) error

GetChunk returns the merged datum hashtrees of a particular chunk (if available)

func (*APIServer) Status

func (a *APIServer) Status(ctx context.Context, _ *types.Empty) (*pps.WorkerStatus, error)

Status returns the status of the current worker.

type CancelRequest

type CancelRequest struct {
	JobID                string   `protobuf:"bytes,2,opt,name=job_id,json=jobId,proto3" json:"job_id,omitempty"`
	DataFilters          []string `protobuf:"bytes,1,rep,name=data_filters,json=dataFilters,proto3" json:"data_filters,omitempty"`
	XXX_NoUnkeyedLiteral struct{} `json:"-"`
	XXX_unrecognized     []byte   `json:"-"`
	XXX_sizecache        int32    `json:"-"`
}

func (*CancelRequest) Descriptor

func (*CancelRequest) Descriptor() ([]byte, []int)

func (*CancelRequest) GetDataFilters

func (m *CancelRequest) GetDataFilters() []string

func (*CancelRequest) GetJobID

func (m *CancelRequest) GetJobID() string

func (*CancelRequest) Marshal

func (m *CancelRequest) Marshal() (dAtA []byte, err error)

func (*CancelRequest) MarshalTo

func (m *CancelRequest) MarshalTo(dAtA []byte) (int, error)

func (*CancelRequest) MarshalToSizedBuffer added in v1.8.8

func (m *CancelRequest) MarshalToSizedBuffer(dAtA []byte) (int, error)

func (*CancelRequest) ProtoMessage

func (*CancelRequest) ProtoMessage()

func (*CancelRequest) Reset

func (m *CancelRequest) Reset()

func (*CancelRequest) Size

func (m *CancelRequest) Size() (n int)

func (*CancelRequest) String

func (m *CancelRequest) String() string

func (*CancelRequest) Unmarshal

func (m *CancelRequest) Unmarshal(dAtA []byte) error

func (*CancelRequest) XXX_DiscardUnknown added in v1.7.12

func (m *CancelRequest) XXX_DiscardUnknown()

func (*CancelRequest) XXX_Marshal added in v1.7.12

func (m *CancelRequest) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*CancelRequest) XXX_Merge added in v1.7.12

func (m *CancelRequest) XXX_Merge(src proto.Message)

func (*CancelRequest) XXX_Size added in v1.7.12

func (m *CancelRequest) XXX_Size() int

func (*CancelRequest) XXX_Unmarshal added in v1.7.12

func (m *CancelRequest) XXX_Unmarshal(b []byte) error

type CancelResponse

type CancelResponse struct {
	Success              bool     `protobuf:"varint,1,opt,name=success,proto3" json:"success,omitempty"`
	XXX_NoUnkeyedLiteral struct{} `json:"-"`
	XXX_unrecognized     []byte   `json:"-"`
	XXX_sizecache        int32    `json:"-"`
}

func (*CancelResponse) Descriptor

func (*CancelResponse) Descriptor() ([]byte, []int)

func (*CancelResponse) GetSuccess

func (m *CancelResponse) GetSuccess() bool

func (*CancelResponse) Marshal

func (m *CancelResponse) Marshal() (dAtA []byte, err error)

func (*CancelResponse) MarshalTo

func (m *CancelResponse) MarshalTo(dAtA []byte) (int, error)

func (*CancelResponse) MarshalToSizedBuffer added in v1.8.8

func (m *CancelResponse) MarshalToSizedBuffer(dAtA []byte) (int, error)

func (*CancelResponse) ProtoMessage

func (*CancelResponse) ProtoMessage()

func (*CancelResponse) Reset

func (m *CancelResponse) Reset()

func (*CancelResponse) Size

func (m *CancelResponse) Size() (n int)

func (*CancelResponse) String

func (m *CancelResponse) String() string

func (*CancelResponse) Unmarshal

func (m *CancelResponse) Unmarshal(dAtA []byte) error

func (*CancelResponse) XXX_DiscardUnknown added in v1.7.12

func (m *CancelResponse) XXX_DiscardUnknown()

func (*CancelResponse) XXX_Marshal added in v1.7.12

func (m *CancelResponse) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*CancelResponse) XXX_Merge added in v1.7.12

func (m *CancelResponse) XXX_Merge(src proto.Message)

func (*CancelResponse) XXX_Size added in v1.7.12

func (m *CancelResponse) XXX_Size() int

func (*CancelResponse) XXX_Unmarshal added in v1.7.12

func (m *CancelResponse) XXX_Unmarshal(b []byte) error

type ChunkState added in v1.6.6

type ChunkState struct {
	State   State  `protobuf:"varint,1,opt,name=state,proto3,enum=worker.State" json:"state,omitempty"`
	DatumID string `protobuf:"bytes,2,opt,name=datum_id,json=datumId,proto3" json:"datum_id,omitempty"`
	// The IP address of the worker who processed this chunk
	Address              string      `protobuf:"bytes,3,opt,name=address,proto3" json:"address,omitempty"`
	RecoveredDatums      *pfs.Object `protobuf:"bytes,4,opt,name=recovered_datums,json=recoveredDatums,proto3" json:"recovered_datums,omitempty"`
	XXX_NoUnkeyedLiteral struct{}    `json:"-"`
	XXX_unrecognized     []byte      `json:"-"`
	XXX_sizecache        int32       `json:"-"`
}

func (*ChunkState) Descriptor added in v1.6.6

func (*ChunkState) Descriptor() ([]byte, []int)

func (*ChunkState) GetAddress added in v1.8.5

func (m *ChunkState) GetAddress() string

func (*ChunkState) GetDatumID added in v1.6.6

func (m *ChunkState) GetDatumID() string

func (*ChunkState) GetRecoveredDatums added in v1.9.9

func (m *ChunkState) GetRecoveredDatums() *pfs.Object

func (*ChunkState) GetState added in v1.6.6

func (m *ChunkState) GetState() State

func (*ChunkState) Marshal added in v1.6.6

func (m *ChunkState) Marshal() (dAtA []byte, err error)

func (*ChunkState) MarshalTo added in v1.6.6

func (m *ChunkState) MarshalTo(dAtA []byte) (int, error)

func (*ChunkState) MarshalToSizedBuffer added in v1.8.8

func (m *ChunkState) MarshalToSizedBuffer(dAtA []byte) (int, error)

func (*ChunkState) ProtoMessage added in v1.6.6

func (*ChunkState) ProtoMessage()

func (*ChunkState) Reset added in v1.6.6

func (m *ChunkState) Reset()

func (*ChunkState) Size added in v1.6.6

func (m *ChunkState) Size() (n int)

func (*ChunkState) String added in v1.6.6

func (m *ChunkState) String() string

func (*ChunkState) Unmarshal added in v1.6.6

func (m *ChunkState) Unmarshal(dAtA []byte) error

func (*ChunkState) XXX_DiscardUnknown added in v1.7.12

func (m *ChunkState) XXX_DiscardUnknown()

func (*ChunkState) XXX_Marshal added in v1.7.12

func (m *ChunkState) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*ChunkState) XXX_Merge added in v1.7.12

func (m *ChunkState) XXX_Merge(src proto.Message)

func (*ChunkState) XXX_Size added in v1.7.12

func (m *ChunkState) XXX_Size() int

func (*ChunkState) XXX_Unmarshal added in v1.7.12

func (m *ChunkState) XXX_Unmarshal(b []byte) error

type Client added in v1.7.5

type Client struct {
	WorkerClient
	debug.DebugClient
}

Client combines the WorkerAPI and the DebugAPI into a single client.

func Clients added in v1.7.5

func Clients(ctx context.Context, pipelineRcName string, etcdClient *etcd.Client, etcdPrefix string, workerGrpcPort uint16) ([]Client, error)

Clients returns a slice of worker clients for a pipeline. pipelineRcName is the name of the pipeline's RC and can be gotten with ppsutil.PipelineRcName. You can also pass "" for pipelineRcName to get all clients for all workers.

func NewClient added in v1.8.5

func NewClient(address string) (Client, error)

NewClient returns a worker client for the worker at the IP address passed in.

type DatumIterator added in v1.8.8

type DatumIterator interface {
	Reset()
	Len() int
	Next() bool
	Datum() []*Input
	DatumN(int) []*Input
}

DatumIterator is an interface which allows you to iterate through the datums for a job. A datum iterator keeps track of which datum it is on, which can be Reset() The intended use is by using this pattern `for di.Next() { ... datum := di.Datum() ... }` Note that since you start the loop by a call to Next(), the datum iterator's location starts at -1

func NewDatumIterator added in v1.8.8

func NewDatumIterator(pachClient *client.APIClient, input *pps.Input) (DatumIterator, error)

NewDatumIterator creates a datumIterator for an input.

type GetChunkRequest added in v1.8.5

type GetChunkRequest struct {
	Id                   int64    `protobuf:"varint,1,opt,name=id,proto3" json:"id,omitempty"`
	Shard                int64    `protobuf:"varint,2,opt,name=shard,proto3" json:"shard,omitempty"`
	Stats                bool     `protobuf:"varint,3,opt,name=stats,proto3" json:"stats,omitempty"`
	XXX_NoUnkeyedLiteral struct{} `json:"-"`
	XXX_unrecognized     []byte   `json:"-"`
	XXX_sizecache        int32    `json:"-"`
}

func (*GetChunkRequest) Descriptor added in v1.8.5

func (*GetChunkRequest) Descriptor() ([]byte, []int)

func (*GetChunkRequest) GetId added in v1.8.5

func (m *GetChunkRequest) GetId() int64

func (*GetChunkRequest) GetShard added in v1.8.5

func (m *GetChunkRequest) GetShard() int64

func (*GetChunkRequest) GetStats added in v1.8.5

func (m *GetChunkRequest) GetStats() bool

func (*GetChunkRequest) Marshal added in v1.8.5

func (m *GetChunkRequest) Marshal() (dAtA []byte, err error)

func (*GetChunkRequest) MarshalTo added in v1.8.5

func (m *GetChunkRequest) MarshalTo(dAtA []byte) (int, error)

func (*GetChunkRequest) MarshalToSizedBuffer added in v1.8.8

func (m *GetChunkRequest) MarshalToSizedBuffer(dAtA []byte) (int, error)

func (*GetChunkRequest) ProtoMessage added in v1.8.5

func (*GetChunkRequest) ProtoMessage()

func (*GetChunkRequest) Reset added in v1.8.5

func (m *GetChunkRequest) Reset()

func (*GetChunkRequest) Size added in v1.8.5

func (m *GetChunkRequest) Size() (n int)

func (*GetChunkRequest) String added in v1.8.5

func (m *GetChunkRequest) String() string

func (*GetChunkRequest) Unmarshal added in v1.8.5

func (m *GetChunkRequest) Unmarshal(dAtA []byte) error

func (*GetChunkRequest) XXX_DiscardUnknown added in v1.8.5

func (m *GetChunkRequest) XXX_DiscardUnknown()

func (*GetChunkRequest) XXX_Marshal added in v1.8.5

func (m *GetChunkRequest) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*GetChunkRequest) XXX_Merge added in v1.8.5

func (m *GetChunkRequest) XXX_Merge(src proto.Message)

func (*GetChunkRequest) XXX_Size added in v1.8.5

func (m *GetChunkRequest) XXX_Size() int

func (*GetChunkRequest) XXX_Unmarshal added in v1.8.5

func (m *GetChunkRequest) XXX_Unmarshal(b []byte) error

type Input

type Input struct {
	FileInfo             *pfs.FileInfo `protobuf:"bytes,1,opt,name=file_info,json=fileInfo,proto3" json:"file_info,omitempty"`
	ParentCommit         *pfs.Commit   `protobuf:"bytes,5,opt,name=parent_commit,json=parentCommit,proto3" json:"parent_commit,omitempty"`
	Name                 string        `protobuf:"bytes,2,opt,name=name,proto3" json:"name,omitempty"`
	JoinOn               string        `protobuf:"bytes,8,opt,name=join_on,json=joinOn,proto3" json:"join_on,omitempty"`
	Lazy                 bool          `protobuf:"varint,3,opt,name=lazy,proto3" json:"lazy,omitempty"`
	Branch               string        `protobuf:"bytes,4,opt,name=branch,proto3" json:"branch,omitempty"`
	GitURL               string        `protobuf:"bytes,6,opt,name=git_url,json=gitUrl,proto3" json:"git_url,omitempty"`
	EmptyFiles           bool          `protobuf:"varint,7,opt,name=empty_files,json=emptyFiles,proto3" json:"empty_files,omitempty"`
	XXX_NoUnkeyedLiteral struct{}      `json:"-"`
	XXX_unrecognized     []byte        `json:"-"`
	XXX_sizecache        int32         `json:"-"`
}

func (*Input) Descriptor

func (*Input) Descriptor() ([]byte, []int)

func (*Input) GetBranch

func (m *Input) GetBranch() string

func (*Input) GetEmptyFiles added in v1.6.8

func (m *Input) GetEmptyFiles() bool

func (*Input) GetFileInfo

func (m *Input) GetFileInfo() *pfs.FileInfo

func (*Input) GetGitURL added in v1.6.4

func (m *Input) GetGitURL() string

func (*Input) GetJoinOn added in v1.8.8

func (m *Input) GetJoinOn() string

func (*Input) GetLazy

func (m *Input) GetLazy() bool

func (*Input) GetName

func (m *Input) GetName() string

func (*Input) GetParentCommit

func (m *Input) GetParentCommit() *pfs.Commit

func (*Input) Marshal

func (m *Input) Marshal() (dAtA []byte, err error)

func (*Input) MarshalTo

func (m *Input) MarshalTo(dAtA []byte) (int, error)

func (*Input) MarshalToSizedBuffer added in v1.8.8

func (m *Input) MarshalToSizedBuffer(dAtA []byte) (int, error)

func (*Input) ProtoMessage

func (*Input) ProtoMessage()

func (*Input) Reset

func (m *Input) Reset()

func (*Input) Size

func (m *Input) Size() (n int)

func (*Input) String

func (m *Input) String() string

func (*Input) Unmarshal

func (m *Input) Unmarshal(dAtA []byte) error

func (*Input) XXX_DiscardUnknown added in v1.7.12

func (m *Input) XXX_DiscardUnknown()

func (*Input) XXX_Marshal added in v1.7.12

func (m *Input) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*Input) XXX_Merge added in v1.7.12

func (m *Input) XXX_Merge(src proto.Message)

func (*Input) XXX_Size added in v1.7.12

func (m *Input) XXX_Size() int

func (*Input) XXX_Unmarshal added in v1.7.12

func (m *Input) XXX_Unmarshal(b []byte) error

type MergeState added in v1.8.0

type MergeState struct {
	State                State       `protobuf:"varint,1,opt,name=state,proto3,enum=worker.State" json:"state,omitempty"`
	Tree                 *pfs.Object `protobuf:"bytes,2,opt,name=tree,proto3" json:"tree,omitempty"`
	SizeBytes            uint64      `protobuf:"varint,3,opt,name=size_bytes,json=sizeBytes,proto3" json:"size_bytes,omitempty"`
	StatsTree            *pfs.Object `protobuf:"bytes,4,opt,name=stats_tree,json=statsTree,proto3" json:"stats_tree,omitempty"`
	StatsSizeBytes       uint64      `protobuf:"varint,5,opt,name=stats_size_bytes,json=statsSizeBytes,proto3" json:"stats_size_bytes,omitempty"`
	XXX_NoUnkeyedLiteral struct{}    `json:"-"`
	XXX_unrecognized     []byte      `json:"-"`
	XXX_sizecache        int32       `json:"-"`
}

func (*MergeState) Descriptor added in v1.8.0

func (*MergeState) Descriptor() ([]byte, []int)

func (*MergeState) GetSizeBytes added in v1.8.0

func (m *MergeState) GetSizeBytes() uint64

func (*MergeState) GetState added in v1.8.0

func (m *MergeState) GetState() State

func (*MergeState) GetStatsSizeBytes added in v1.8.0

func (m *MergeState) GetStatsSizeBytes() uint64

func (*MergeState) GetStatsTree added in v1.8.0

func (m *MergeState) GetStatsTree() *pfs.Object

func (*MergeState) GetTree added in v1.8.0

func (m *MergeState) GetTree() *pfs.Object

func (*MergeState) Marshal added in v1.8.0

func (m *MergeState) Marshal() (dAtA []byte, err error)

func (*MergeState) MarshalTo added in v1.8.0

func (m *MergeState) MarshalTo(dAtA []byte) (int, error)

func (*MergeState) MarshalToSizedBuffer added in v1.8.8

func (m *MergeState) MarshalToSizedBuffer(dAtA []byte) (int, error)

func (*MergeState) ProtoMessage added in v1.8.0

func (*MergeState) ProtoMessage()

func (*MergeState) Reset added in v1.8.0

func (m *MergeState) Reset()

func (*MergeState) Size added in v1.8.0

func (m *MergeState) Size() (n int)

func (*MergeState) String added in v1.8.0

func (m *MergeState) String() string

func (*MergeState) Unmarshal added in v1.8.0

func (m *MergeState) Unmarshal(dAtA []byte) error

func (*MergeState) XXX_DiscardUnknown added in v1.8.1

func (m *MergeState) XXX_DiscardUnknown()

func (*MergeState) XXX_Marshal added in v1.8.1

func (m *MergeState) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*MergeState) XXX_Merge added in v1.8.1

func (m *MergeState) XXX_Merge(src proto.Message)

func (*MergeState) XXX_Size added in v1.8.1

func (m *MergeState) XXX_Size() int

func (*MergeState) XXX_Unmarshal added in v1.8.1

func (m *MergeState) XXX_Unmarshal(b []byte) error

type Plan added in v1.8.0

type Plan struct {
	Chunks               []int64  `protobuf:"varint,1,rep,packed,name=chunks,proto3" json:"chunks,omitempty"`
	Merges               int64    `protobuf:"varint,2,opt,name=merges,proto3" json:"merges,omitempty"`
	XXX_NoUnkeyedLiteral struct{} `json:"-"`
	XXX_unrecognized     []byte   `json:"-"`
	XXX_sizecache        int32    `json:"-"`
}

func (*Plan) Descriptor added in v1.8.0

func (*Plan) Descriptor() ([]byte, []int)

func (*Plan) GetChunks added in v1.8.0

func (m *Plan) GetChunks() []int64

func (*Plan) GetMerges added in v1.8.0

func (m *Plan) GetMerges() int64

func (*Plan) Marshal added in v1.8.0

func (m *Plan) Marshal() (dAtA []byte, err error)

func (*Plan) MarshalTo added in v1.8.0

func (m *Plan) MarshalTo(dAtA []byte) (int, error)

func (*Plan) MarshalToSizedBuffer added in v1.8.8

func (m *Plan) MarshalToSizedBuffer(dAtA []byte) (int, error)

func (*Plan) ProtoMessage added in v1.8.0

func (*Plan) ProtoMessage()

func (*Plan) Reset added in v1.8.0

func (m *Plan) Reset()

func (*Plan) Size added in v1.8.0

func (m *Plan) Size() (n int)

func (*Plan) String added in v1.8.0

func (m *Plan) String() string

func (*Plan) Unmarshal added in v1.8.0

func (m *Plan) Unmarshal(dAtA []byte) error

func (*Plan) XXX_DiscardUnknown added in v1.8.1

func (m *Plan) XXX_DiscardUnknown()

func (*Plan) XXX_Marshal added in v1.8.1

func (m *Plan) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*Plan) XXX_Merge added in v1.8.1

func (m *Plan) XXX_Merge(src proto.Message)

func (*Plan) XXX_Size added in v1.8.1

func (m *Plan) XXX_Size() int

func (*Plan) XXX_Unmarshal added in v1.8.1

func (m *Plan) XXX_Unmarshal(b []byte) error

type ShardInfo added in v1.8.5

type ShardInfo struct {
	XXX_NoUnkeyedLiteral struct{} `json:"-"`
	XXX_unrecognized     []byte   `json:"-"`
	XXX_sizecache        int32    `json:"-"`
}

func (*ShardInfo) Descriptor added in v1.8.5

func (*ShardInfo) Descriptor() ([]byte, []int)

func (*ShardInfo) Marshal added in v1.8.5

func (m *ShardInfo) Marshal() (dAtA []byte, err error)

func (*ShardInfo) MarshalTo added in v1.8.5

func (m *ShardInfo) MarshalTo(dAtA []byte) (int, error)

func (*ShardInfo) MarshalToSizedBuffer added in v1.8.8

func (m *ShardInfo) MarshalToSizedBuffer(dAtA []byte) (int, error)

func (*ShardInfo) ProtoMessage added in v1.8.5

func (*ShardInfo) ProtoMessage()

func (*ShardInfo) Reset added in v1.8.5

func (m *ShardInfo) Reset()

func (*ShardInfo) Size added in v1.8.5

func (m *ShardInfo) Size() (n int)

func (*ShardInfo) String added in v1.8.5

func (m *ShardInfo) String() string

func (*ShardInfo) Unmarshal added in v1.8.5

func (m *ShardInfo) Unmarshal(dAtA []byte) error

func (*ShardInfo) XXX_DiscardUnknown added in v1.8.5

func (m *ShardInfo) XXX_DiscardUnknown()

func (*ShardInfo) XXX_Marshal added in v1.8.5

func (m *ShardInfo) XXX_Marshal(b []byte, deterministic bool) ([]byte, error)

func (*ShardInfo) XXX_Merge added in v1.8.5

func (m *ShardInfo) XXX_Merge(src proto.Message)

func (*ShardInfo) XXX_Size added in v1.8.5

func (m *ShardInfo) XXX_Size() int

func (*ShardInfo) XXX_Unmarshal added in v1.8.5

func (m *ShardInfo) XXX_Unmarshal(b []byte) error

type State added in v1.8.0

type State int32
const (
	State_RUNNING  State = 0
	State_COMPLETE State = 1
	State_FAILED   State = 3
)

func (State) EnumDescriptor added in v1.8.0

func (State) EnumDescriptor() ([]byte, []int)

func (State) String added in v1.8.0

func (x State) String() string

type UnimplementedWorkerServer added in v1.8.7

type UnimplementedWorkerServer struct {
}

UnimplementedWorkerServer can be embedded to have forward compatible implementations.

func (*UnimplementedWorkerServer) Cancel added in v1.8.7

func (*UnimplementedWorkerServer) GetChunk added in v1.8.7

func (*UnimplementedWorkerServer) Status added in v1.8.7

type WorkerClient

type WorkerClient interface {
	Status(ctx context.Context, in *types.Empty, opts ...grpc.CallOption) (*pps.WorkerStatus, error)
	Cancel(ctx context.Context, in *CancelRequest, opts ...grpc.CallOption) (*CancelResponse, error)
	GetChunk(ctx context.Context, in *GetChunkRequest, opts ...grpc.CallOption) (Worker_GetChunkClient, error)
}

WorkerClient is the client API for Worker service.

For semantics around ctx use and closing/ending streaming RPCs, please refer to https://godoc.org/google.golang.org/grpc#ClientConn.NewStream.

func NewWorkerClient

func NewWorkerClient(cc *grpc.ClientConn) WorkerClient

type WorkerServer

type WorkerServer interface {
	Status(context.Context, *types.Empty) (*pps.WorkerStatus, error)
	Cancel(context.Context, *CancelRequest) (*CancelResponse, error)
	GetChunk(*GetChunkRequest, Worker_GetChunkServer) error
}

WorkerServer is the server API for Worker service.

type Worker_GetChunkClient added in v1.8.5

type Worker_GetChunkClient interface {
	Recv() (*types.BytesValue, error)
	grpc.ClientStream
}

type Worker_GetChunkServer added in v1.8.5

type Worker_GetChunkServer interface {
	Send(*types.BytesValue) error
	grpc.ServerStream
}

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL