Documentation
¶
Overview ¶
Package driver is the handwritten Go contract for proto/cloudpath/v1/driver.proto (Driver Protocol v1).
Every type maps 1:1 to a proto message and is JSON-encoded exactly the way the proto text defines it (proto3 json_name, snake_case keys, bytes as base64). oneof fields are modeled as a typed Union plus custom MarshalJSON/UnmarshalJSON in codec.go. There is no protoc-generated code, and no third-party dependency is used.
The package also provides DriverClient/DriverServer interfaces and in-process adapters that run the protocol over sdk/go/transport:
cli := driver.NewClient(transportEnd) srv := driver.NewRPCServer(transportEnd, myDriverImpl) go srv.Serve(ctx)
See testing/plugin-harness for a mock Core, mock Driver and the conformance suite.
Index ¶
- Constants
- func NegotiateProtocolVersion(supported []uint32, prefer uint32, minSupported, maxSupported uint32) (uint32, bool)
- func NewRPCServer(tr transport.Transport, impl DriverServer) *rpc.Server
- func ValidateDriverMessage(m *DriverMessage) error
- type ActionDescriptor
- type CapabilityDescriptor
- type CloseDeviceRequest
- type CloseDeviceResponse
- type CommandProgress
- type CommandState
- type ConfigureInstanceRequest
- type ConfigureInstanceResponse
- type DescribeRequest
- type Device
- type DeviceStatus
- type DeviceUpsert
- type Diagnostic
- type DiscoverRequest
- type DiscoveryEvent
- type DiscoveryEventUnion
- type DiscoveryFailed
- type DiscoveryFinished
- type DiscoveryFoundDevice
- type DiscoveryProgress
- type DiscoveryStarted
- type DiscoveryStream
- type DiscoveryWriter
- type DriverClient
- type DriverDescriptor
- type DriverMessage
- type DriverMessageStream
- type DriverMessageUnion
- type DriverMessageWriter
- type DriverServer
- type Entity
- type EntityCategory
- type EntityUpsert
- type Event
- type EventDescriptor
- type ExecuteRequest
- type ExecuteResponse
- type HealthRequest
- type HealthResponse
- type HealthState
- type InitializeRequest
- type InitializeResponse
- type InstanceHealth
- type Observation
- type OpenDeviceRequest
- type OpenDeviceResponse
- type PropertyDescriptor
- type SequenceTracker
- type ShutdownRequest
- type ShutdownResponse
- type Status
- type Value
- type ValueKind
- type WatchRequest
Constants ¶
const ( MethodInitialize = "cloudpath.v1.driver.DriverService/Initialize" MethodDescribe = "cloudpath.v1.driver.DriverService/Describe" MethodConfigureInstance = "cloudpath.v1.driver.DriverService/ConfigureInstance" MethodDiscover = "cloudpath.v1.driver.DriverService/Discover" MethodOpenDevice = "cloudpath.v1.driver.DriverService/OpenDevice" MethodCloseDevice = "cloudpath.v1.driver.DriverService/CloseDevice" MethodWatch = "cloudpath.v1.driver.DriverService/Watch" MethodExecute = "cloudpath.v1.driver.DriverService/Execute" MethodHealth = "cloudpath.v1.driver.DriverService/Health" MethodShutdown = "cloudpath.v1.driver.DriverService/Shutdown" )
Service method names carried on the wire.
const ProtocolVersion uint32 = 1
ProtocolVersion is the single driver protocol version implemented here.
const SchemaVersion = "1"
SchemaVersion is the data semantic version stamped on every stream message.
Variables ¶
This section is empty.
Functions ¶
func NegotiateProtocolVersion ¶
func NegotiateProtocolVersion(supported []uint32, prefer uint32, minSupported, maxSupported uint32) (uint32, bool)
NegotiateProtocolVersion picks the highest version that is supported by both sides and lies inside [minSupported, maxSupported]. It returns (0, false) when the sets are disjoint. prefer is the caller's preferred version and breaks ties when both sides support multiple versions.
func NewRPCServer ¶
func NewRPCServer(tr transport.Transport, impl DriverServer) *rpc.Server
NewRPCServer binds impl to a transport and returns the dispatcher. Call Serve to consume frames.
func ValidateDriverMessage ¶
func ValidateDriverMessage(m *DriverMessage) error
ValidateDriverMessage checks the report envelope invariants that every implementation must enforce before sending: instance ID, positive sequence, schema version and a body. It returns a descriptive error.
Types ¶
type ActionDescriptor ¶
type ActionDescriptor struct {
Name string `json:"name"`
InputSchemaJSON string `json:"input_schema_json"`
ResultSchemaJSON string `json:"result_schema_json"`
// Optional Capability metadata; omitted values preserve legacy descriptors.
Title string `json:"title,omitempty"`
Description string `json:"description,omitempty"`
Destructive bool `json:"destructive,omitempty"`
Confirmation string `json:"confirmation,omitempty"`
}
type CapabilityDescriptor ¶
type CapabilityDescriptor struct {
ID string `json:"id"`
Title string `json:"title"`
Properties []PropertyDescriptor `json:"properties"`
Events []EventDescriptor `json:"events"`
Actions []ActionDescriptor `json:"actions"`
}
type CloseDeviceRequest ¶
type CloseDeviceResponse ¶
type CommandProgress ¶
type CommandProgress struct {
CommandID string `json:"command_id"`
IdempotencyKey string `json:"idempotency_key"`
EntityID string `json:"entity_id"`
Action string `json:"action"`
State CommandState `json:"state"`
Progress float64 `json:"progress"`
Detail string `json:"detail"`
ResultJSON string `json:"result_json"`
ErrorCode string `json:"error_code"`
}
type CommandState ¶
type CommandState int32
const ( CommandStateUnspecified CommandState = 0 CommandStateCreated CommandState = 1 CommandStateDispatched CommandState = 2 CommandStateAccepted CommandState = 3 CommandStateRunning CommandState = 4 CommandStateSucceeded CommandState = 5 CommandStateFailed CommandState = 6 CommandStateTimedOut CommandState = 7 CommandStateCancelled CommandState = 8 )
func (CommandState) String ¶
func (s CommandState) String() string
String returns a stable, human-readable name for the state.
type DescribeRequest ¶
type DescribeRequest struct{}
DescribeRequest is empty by contract; the descriptor must be stable and deterministic for a given plugin version.
type DeviceStatus ¶
type DeviceStatus int32
const ( DeviceStatusUnspecified DeviceStatus = 0 DeviceStatusOnline DeviceStatus = 1 DeviceStatusOffline DeviceStatus = 2 DeviceStatusDegraded DeviceStatus = 4 )
func (DeviceStatus) String ¶
func (s DeviceStatus) String() string
String returns a stable, human-readable name for the device status.
type DeviceUpsert ¶
type Diagnostic ¶
type DiscoverRequest ¶
type DiscoveryEvent ¶
type DiscoveryEvent struct {
PluginInstanceID string `json:"plugin_instance_id"`
Sequence uint64 `json:"sequence"`
SchemaVersion string `json:"schema_version"`
DiscoveryID string `json:"discovery_id"`
Union DiscoveryEventUnion `json:"-"`
}
DiscoveryEvent carries exactly one body variant.
func (*DiscoveryEvent) MarshalJSON ¶
func (e *DiscoveryEvent) MarshalJSON() ([]byte, error)
MarshalJSON implements the DiscoveryEvent oneof.
func (*DiscoveryEvent) UnmarshalJSON ¶
func (e *DiscoveryEvent) UnmarshalJSON(data []byte) error
UnmarshalJSON implements the DiscoveryEvent oneof.
type DiscoveryEventUnion ¶
type DiscoveryEventUnion interface {
// contains filtered or unexported methods
}
type DiscoveryFailed ¶
type DiscoveryFailed struct {
Status *Status `json:"status"`
}
type DiscoveryFinished ¶
type DiscoveryFinished struct {
FoundCount uint32 `json:"found_count"`
}
type DiscoveryFoundDevice ¶
type DiscoveryProgress ¶
type DiscoveryStarted ¶
type DiscoveryStarted struct {
DiscoveryID string `json:"discovery_id"`
}
type DiscoveryStream ¶
type DiscoveryStream interface {
Recv(ctx context.Context) (*DiscoveryEvent, error)
Cancel(ctx context.Context) error
}
DiscoveryStream is the client view of the Discover server stream.
type DiscoveryWriter ¶
type DiscoveryWriter interface {
Send(ctx context.Context, msg *DiscoveryEvent) error
}
DiscoveryWriter lets a plugin publish discovery events.
type DriverClient ¶
type DriverClient interface {
Initialize(ctx context.Context, req *InitializeRequest) (*InitializeResponse, error)
Describe(ctx context.Context) (*DriverDescriptor, error)
ConfigureInstance(ctx context.Context, req *ConfigureInstanceRequest) (*ConfigureInstanceResponse, error)
Discover(ctx context.Context, req *DiscoverRequest) (DiscoveryStream, error)
OpenDevice(ctx context.Context, req *OpenDeviceRequest) (*OpenDeviceResponse, error)
CloseDevice(ctx context.Context, req *CloseDeviceRequest) (*CloseDeviceResponse, error)
Watch(ctx context.Context, req *WatchRequest) (DriverMessageStream, error)
Execute(ctx context.Context, req *ExecuteRequest) (*ExecuteResponse, error)
Health(ctx context.Context) (*HealthResponse, error)
Shutdown(ctx context.Context, req *ShutdownRequest) (*ShutdownResponse, error)
}
DriverClient is the Core-side view of DriverService.
func NewClient ¶
func NewClient(tr transport.Transport) DriverClient
NewClient wraps a transport as a DriverClient.
type DriverDescriptor ¶
type DriverDescriptor struct {
DriverID string `json:"driver_id"`
Version string `json:"version"`
SchemaVersions []string `json:"schema_versions"`
Capabilities []CapabilityDescriptor `json:"capabilities"`
}
type DriverMessage ¶
type DriverMessage struct {
PluginInstanceID string `json:"plugin_instance_id"`
Sequence uint64 `json:"sequence"`
SchemaVersion string `json:"schema_version"`
DeviceID string `json:"device_id"`
Union DriverMessageUnion `json:"-"`
}
DriverMessage is the single report envelope of the Watch stream. Every message carries plugin_instance_id, a monotonically increasing sequence scoped to (instance, device), and a schema version. Core dedupes on that scope; see SequenceTracker.
func SortMessagesBySequence ¶
func SortMessagesBySequence(msgs []*DriverMessage) []*DriverMessage
SortMessagesBySequence is a helper for replay buffers: it returns a copy sorted ascending by sequence (stable) without mutating the input.
func (*DriverMessage) MarshalJSON ¶
func (m *DriverMessage) MarshalJSON() ([]byte, error)
MarshalJSON implements the DriverMessage oneof.
func (*DriverMessage) UnmarshalJSON ¶
func (m *DriverMessage) UnmarshalJSON(data []byte) error
UnmarshalJSON implements the DriverMessage oneof and rejects messages with zero or multiple bodies.
type DriverMessageStream ¶
type DriverMessageStream interface {
Recv(ctx context.Context) (*DriverMessage, error)
Cancel(ctx context.Context) error
}
DriverMessageStream is the client view of the Watch server stream.
type DriverMessageUnion ¶
type DriverMessageUnion interface {
// contains filtered or unexported methods
}
type DriverMessageWriter ¶
type DriverMessageWriter interface {
Send(ctx context.Context, msg *DriverMessage) error
}
DriverMessageWriter lets a plugin publish Watch messages.
type DriverServer ¶
type DriverServer interface {
Initialize(ctx context.Context, req *InitializeRequest) (*InitializeResponse, error)
Describe(ctx context.Context) (*DriverDescriptor, error)
ConfigureInstance(ctx context.Context, req *ConfigureInstanceRequest) (*ConfigureInstanceResponse, error)
Discover(ctx context.Context, req *DiscoverRequest, stream DiscoveryWriter) error
OpenDevice(ctx context.Context, req *OpenDeviceRequest) (*OpenDeviceResponse, error)
CloseDevice(ctx context.Context, req *CloseDeviceRequest) (*CloseDeviceResponse, error)
Watch(ctx context.Context, req *WatchRequest, stream DriverMessageWriter) error
Execute(ctx context.Context, req *ExecuteRequest) (*ExecuteResponse, error)
Health(ctx context.Context) (*HealthResponse, error)
Shutdown(ctx context.Context, req *ShutdownRequest) (*ShutdownResponse, error)
}
DriverServer is the plugin-side implementation of DriverService.
type EntityCategory ¶
type EntityCategory int32
const ( EntityCategoryUnspecified EntityCategory = 0 EntityCategorySensor EntityCategory = 1 EntityCategoryActuator EntityCategory = 2 EntityCategoryDiagnostic EntityCategory = 3 EntityCategoryConfig EntityCategory = 4 )
type EntityUpsert ¶
type EventDescriptor ¶
type ExecuteRequest ¶
type ExecuteRequest struct {
PluginInstanceID string `json:"plugin_instance_id"`
DeviceID string `json:"device_id,omitempty"`
IdempotencyKey string `json:"idempotency_key"`
EntityID string `json:"entity_id"`
Action string `json:"action"`
ArgsJSON string `json:"args_json"`
Deadline string `json:"deadline"`
Actor string `json:"actor"`
CancelCommandID string `json:"cancel_command_id"`
}
ExecuteRequest either runs a command or, when CancelCommandID is non-empty, cancels the referenced in-flight command. IdempotencyKey is required in both cases so Core can safely retry.
type ExecuteResponse ¶
type HealthRequest ¶
type HealthRequest struct{}
type HealthResponse ¶
type HealthResponse struct {
State HealthState `json:"state"`
Instances []InstanceHealth `json:"instances"`
}
type HealthState ¶
type HealthState int32
const ( HealthStateUnspecified HealthState = 0 HealthStateServing HealthState = 1 HealthStateNotServing HealthState = 2 )
type InitializeRequest ¶
type InitializeRequest struct {
PluginID string `json:"plugin_id"`
PluginVersion string `json:"plugin_version"`
LaunchID string `json:"launch_id"`
HandshakeCookie string `json:"handshake_cookie"`
ProtocolVersion uint32 `json:"protocol_version"`
SupportedProtocolVersions []uint32 `json:"supported_protocol_versions"`
NodeID string `json:"node_id"`
RuntimeType string `json:"runtime_type"`
HostInfo map[string]string `json:"host_info"`
}
type InitializeResponse ¶
type InstanceHealth ¶
type InstanceHealth struct {
PluginInstanceID string `json:"plugin_instance_id"`
State HealthState `json:"state"`
Detail string `json:"detail"`
}
type Observation ¶
type OpenDeviceRequest ¶
type OpenDeviceResponse ¶
type PropertyDescriptor ¶
type SequenceTracker ¶
type SequenceTracker struct {
// contains filtered or unexported fields
}
SequenceTracker implements Core-side deduplication: for a given (plugin_instance_id, device_id) scope only a sequence strictly greater than the last accepted one is admitted. Duplicates and stale/out-of-order frames are therefore dropped at the boundary.
func NewSequenceTracker ¶
func NewSequenceTracker() *SequenceTracker
NewSequenceTracker returns an empty tracker.
func (*SequenceTracker) Accept ¶
func (t *SequenceTracker) Accept(instanceID, deviceID string, sequence uint64) bool
Accept reports whether sequence is new for the scope and, when it is, records it. Sequence 0 is rejected because protocol messages must carry a positive sequence.
func (*SequenceTracker) AcceptMessage ¶
func (t *SequenceTracker) AcceptMessage(m *DriverMessage) bool
AcceptMessage applies Accept to a DriverMessage's own scope fields.
func (*SequenceTracker) Count ¶
func (t *SequenceTracker) Count() int
Count returns how many distinct scopes have been observed.
func (*SequenceTracker) Last ¶
func (t *SequenceTracker) Last(instanceID, deviceID string) uint64
Last returns the highest accepted sequence for the scope (0 if none).
func (*SequenceTracker) Reset ¶
func (t *SequenceTracker) Reset()
Reset clears all scopes. Useful when a plugin restarts and replays from a fresh sequence space.
type ShutdownRequest ¶
type ShutdownResponse ¶
type ShutdownResponse struct {
Status *Status `json:"status"`
}
type Status ¶
Status is the protocol status type (mirrors the inline Status message in driver.proto; the application protocol shares the same code table).
type Value ¶
type Value struct {
Kind ValueKind `json:"-"`
NumberValue float64 `json:"-"`
IntValue int64 `json:"-"`
StringValue string `json:"-"`
BoolValue bool `json:"-"`
JSONValue string `json:"-"`
}
Value is the typed oneof inside Observation.
func (*Value) MarshalJSON ¶
MarshalJSON implements the Value oneof.
func (*Value) UnmarshalJSON ¶
UnmarshalJSON implements the Value oneof.