From 3ea2e9eb7220b11fbdfc2e155f35ba8ba6cb4c33 Mon Sep 17 00:00:00 2001 From: dawn Date: Mon, 6 Jul 2026 19:26:57 +0300 Subject: [PATCH] spindle/mill: add executor wire protocol --- buf.gen.yaml | 4 + buf.yaml | 1 + spindle/mill/proto/gen/mill.pb.go | 1317 +++++++++++++++++ spindle/mill/proto/protocol.go | 103 ++ spindle/mill/proto/protocol_test.go | 66 + spindle/mill/proto/spindle/mill/v1/mill.proto | 169 +++ spindle/mill/proto/ws.go | 58 + 7 files changed, 1718 insertions(+) create mode 100644 spindle/mill/proto/gen/mill.pb.go create mode 100644 spindle/mill/proto/protocol.go create mode 100644 spindle/mill/proto/protocol_test.go create mode 100644 spindle/mill/proto/spindle/mill/v1/mill.proto create mode 100644 spindle/mill/proto/ws.go diff --git a/buf.gen.yaml b/buf.gen.yaml index 7286491c..7681c279 100644 --- a/buf.gen.yaml +++ b/buf.gen.yaml @@ -9,3 +9,7 @@ plugins: out: shuttle/src/gen opt: - bytes=. + # the fleet protocol is broker<->executor only (both Go); shuttle (Rust) + # only speaks the agent protocol, so keep fleet types out of its gen tree. + exclude_types: + - spindle.mill.v1 diff --git a/buf.yaml b/buf.yaml index 3c20343f..7b664379 100644 --- a/buf.yaml +++ b/buf.yaml @@ -1,5 +1,6 @@ version: v2 modules: - path: spindle/agentproto + - path: spindle/mill/proto deps: - buf.build/bufbuild/protovalidate diff --git a/spindle/mill/proto/gen/mill.pb.go b/spindle/mill/proto/gen/mill.pb.go new file mode 100644 index 00000000..8e539a4b --- /dev/null +++ b/spindle/mill/proto/gen/mill.pb.go @@ -0,0 +1,1317 @@ +// Code generated by protoc-gen-go. DO NOT EDIT. +// versions: +// protoc-gen-go v1.36.11 +// protoc (unknown) +// source: spindle/mill/v1/mill.proto + +package millv1 + +import ( + _ "buf.build/gen/go/bufbuild/protovalidate/protocolbuffers/go/buf/validate" + protoreflect "google.golang.org/protobuf/reflect/protoreflect" + protoimpl "google.golang.org/protobuf/runtime/protoimpl" + reflect "reflect" + sync "sync" + unsafe "unsafe" +) + +const ( + // Verify that this generated code is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(20 - protoimpl.MinVersion) + // Verify that runtime/protoimpl is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) +) + +// Hello is the first frame an executor sends after dialing the mill. It +// carries the static traits of the node plus a resume hint. +type Hello struct { + state protoimpl.MessageState `protogen:"open.v1"` + ProtocolVersion uint32 `protobuf:"varint,1,opt,name=protocol_version,json=protocolVersion,proto3" json:"protocol_version,omitempty"` + // stable across reconnects (e.g. hostname / DID); the mill keys sessions on + // this so a brief blip reattaches the same node rather than creating a new one. + NodeId string `protobuf:"bytes,2,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` + // engine names this node can run ("microvm", "nixery"). + Engines []string `protobuf:"bytes,3,rep,name=engines,proto3" json:"engines,omitempty"` + // GOARCH of the node, so the mill won't place arch-incompatible jobs. + Arch string `protobuf:"bytes,4,opt,name=arch,proto3" json:"arch,omitempty"` + // the highest relay offset the executor believes it has sent; a resume hint. + LastOffset uint64 `protobuf:"varint,5,opt,name=last_offset,json=lastOffset,proto3" json:"last_offset,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *Hello) Reset() { + *x = Hello{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[0] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Hello) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Hello) ProtoMessage() {} + +func (x *Hello) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[0] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Hello.ProtoReflect.Descriptor instead. +func (*Hello) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{0} +} + +func (x *Hello) GetProtocolVersion() uint32 { + if x != nil { + return x.ProtocolVersion + } + return 0 +} + +func (x *Hello) GetNodeId() string { + if x != nil { + return x.NodeId + } + return "" +} + +func (x *Hello) GetEngines() []string { + if x != nil { + return x.Engines + } + return nil +} + +func (x *Hello) GetArch() string { + if x != nil { + return x.Arch + } + return "" +} + +func (x *Hello) GetLastOffset() uint64 { + if x != nil { + return x.LastOffset + } + return 0 +} + +// Resume is the mill's reply to Hello. The executor replays every buffered +// relay entry with offset strictly greater than ack_offset before sending new +// ones. +type Resume struct { + state protoimpl.MessageState `protogen:"open.v1"` + AckOffset uint64 `protobuf:"varint,1,opt,name=ack_offset,json=ackOffset,proto3" json:"ack_offset,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *Resume) Reset() { + *x = Resume{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[1] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Resume) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Resume) ProtoMessage() {} + +func (x *Resume) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[1] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Resume.ProtoReflect.Descriptor instead. +func (*Resume) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{1} +} + +func (x *Resume) GetAckOffset() uint64 { + if x != nil { + return x.AckOffset + } + return 0 +} + +// EngineSnapshot is the changing per-engine state on a node. free_seats is a +// coarse "can you take more" hint; 0 also means draining (the mill treats a +// draining node as a full one until it leaves). The resource fields are coarse +// budget headroom for ranking only; the executor makes the real yes/no call in +// ReserveResult. +type EngineSnapshot struct { + state protoimpl.MessageState `protogen:"open.v1"` + FreeSeats uint32 `protobuf:"varint,1,opt,name=free_seats,json=freeSeats,proto3" json:"free_seats,omitempty"` + FreeMemoryMib int64 `protobuf:"varint,2,opt,name=free_memory_mib,json=freeMemoryMib,proto3" json:"free_memory_mib,omitempty"` + FreeVcpus int64 `protobuf:"varint,3,opt,name=free_vcpus,json=freeVcpus,proto3" json:"free_vcpus,omitempty"` + FreeDiskMib int64 `protobuf:"varint,4,opt,name=free_disk_mib,json=freeDiskMib,proto3" json:"free_disk_mib,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *EngineSnapshot) Reset() { + *x = EngineSnapshot{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[2] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *EngineSnapshot) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*EngineSnapshot) ProtoMessage() {} + +func (x *EngineSnapshot) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[2] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use EngineSnapshot.ProtoReflect.Descriptor instead. +func (*EngineSnapshot) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{2} +} + +func (x *EngineSnapshot) GetFreeSeats() uint32 { + if x != nil { + return x.FreeSeats + } + return 0 +} + +func (x *EngineSnapshot) GetFreeMemoryMib() int64 { + if x != nil { + return x.FreeMemoryMib + } + return 0 +} + +func (x *EngineSnapshot) GetFreeVcpus() int64 { + if x != nil { + return x.FreeVcpus + } + return 0 +} + +func (x *EngineSnapshot) GetFreeDiskMib() int64 { + if x != nil { + return x.FreeDiskMib + } + return 0 +} + +// NodeSnapshot is pushed on connect, periodically, and right after any state +// change (reserve, commit, terminal). +type NodeSnapshot struct { + state protoimpl.MessageState `protogen:"open.v1"` + NodeId string `protobuf:"bytes,1,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` + Seq uint64 `protobuf:"varint,2,opt,name=seq,proto3" json:"seq,omitempty"` + Engines map[string]*EngineSnapshot `protobuf:"bytes,3,rep,name=engines,proto3" json:"engines,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *NodeSnapshot) Reset() { + *x = NodeSnapshot{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[3] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *NodeSnapshot) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*NodeSnapshot) ProtoMessage() {} + +func (x *NodeSnapshot) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[3] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use NodeSnapshot.ProtoReflect.Descriptor instead. +func (*NodeSnapshot) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{3} +} + +func (x *NodeSnapshot) GetNodeId() string { + if x != nil { + return x.NodeId + } + return "" +} + +func (x *NodeSnapshot) GetSeq() uint64 { + if x != nil { + return x.Seq + } + return 0 +} + +func (x *NodeSnapshot) GetEngines() map[string]*EngineSnapshot { + if x != nil { + return x.Engines + } + return nil +} + +// ReserveSeat asks an executor to hold a seat for a job. Zero secrets ride this +// message; the raw pipeline/workflow are carried as JSON since processPipeline +// already round-trips them through JSON. +type ReserveSeat struct { + state protoimpl.MessageState `protogen:"open.v1"` + LeaseId string `protobuf:"bytes,1,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + TargetEngine string `protobuf:"bytes,2,opt,name=target_engine,json=targetEngine,proto3" json:"target_engine,omitempty"` + RawPipelineJson string `protobuf:"bytes,3,opt,name=raw_pipeline_json,json=rawPipelineJson,proto3" json:"raw_pipeline_json,omitempty"` + RawWorkflowJson string `protobuf:"bytes,4,opt,name=raw_workflow_json,json=rawWorkflowJson,proto3" json:"raw_workflow_json,omitempty"` + // the pipeline id (knot + rkey); the executor reconstructs the exact + // WorkflowId so its relayed status rows and log path match what the mill + // authored for "pending". + Knot string `protobuf:"bytes,5,opt,name=knot,proto3" json:"knot,omitempty"` + Rkey string `protobuf:"bytes,6,opt,name=rkey,proto3" json:"rkey,omitempty"` + TtlSeconds uint32 `protobuf:"varint,7,opt,name=ttl_seconds,json=ttlSeconds,proto3" json:"ttl_seconds,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ReserveSeat) Reset() { + *x = ReserveSeat{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[4] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ReserveSeat) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ReserveSeat) ProtoMessage() {} + +func (x *ReserveSeat) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[4] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ReserveSeat.ProtoReflect.Descriptor instead. +func (*ReserveSeat) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{4} +} + +func (x *ReserveSeat) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +func (x *ReserveSeat) GetTargetEngine() string { + if x != nil { + return x.TargetEngine + } + return "" +} + +func (x *ReserveSeat) GetRawPipelineJson() string { + if x != nil { + return x.RawPipelineJson + } + return "" +} + +func (x *ReserveSeat) GetRawWorkflowJson() string { + if x != nil { + return x.RawWorkflowJson + } + return "" +} + +func (x *ReserveSeat) GetKnot() string { + if x != nil { + return x.Knot + } + return "" +} + +func (x *ReserveSeat) GetRkey() string { + if x != nil { + return x.Rkey + } + return "" +} + +func (x *ReserveSeat) GetTtlSeconds() uint32 { + if x != nil { + return x.TtlSeconds + } + return 0 +} + +// ReserveResult is the executor's accept/reject for a ReserveSeat. +type ReserveResult struct { + state protoimpl.MessageState `protogen:"open.v1"` + LeaseId string `protobuf:"bytes,1,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + Accepted bool `protobuf:"varint,2,opt,name=accepted,proto3" json:"accepted,omitempty"` + RejectReason string `protobuf:"bytes,3,opt,name=reject_reason,json=rejectReason,proto3" json:"reject_reason,omitempty"` + // optional bid score the mill ranks accepted leases by (higher is better). + Score float64 `protobuf:"fixed64,4,opt,name=score,proto3" json:"score,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ReserveResult) Reset() { + *x = ReserveResult{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[5] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ReserveResult) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ReserveResult) ProtoMessage() {} + +func (x *ReserveResult) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[5] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ReserveResult.ProtoReflect.Descriptor instead. +func (*ReserveResult) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{5} +} + +func (x *ReserveResult) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +func (x *ReserveResult) GetAccepted() bool { + if x != nil { + return x.Accepted + } + return false +} + +func (x *ReserveResult) GetRejectReason() string { + if x != nil { + return x.RejectReason + } + return "" +} + +func (x *ReserveResult) GetScore() float64 { + if x != nil { + return x.Score + } + return 0 +} + +// Secret is a single unlocked secret, sent only inside CommitLease. +type Secret struct { + state protoimpl.MessageState `protogen:"open.v1"` + Key string `protobuf:"bytes,1,opt,name=key,proto3" json:"key,omitempty"` + Value string `protobuf:"bytes,2,opt,name=value,proto3" json:"value,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *Secret) Reset() { + *x = Secret{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[6] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Secret) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Secret) ProtoMessage() {} + +func (x *Secret) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[6] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Secret.ProtoReflect.Descriptor instead. +func (*Secret) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{6} +} + +func (x *Secret) GetKey() string { + if x != nil { + return x.Key + } + return "" +} + +func (x *Secret) GetValue() string { + if x != nil { + return x.Value + } + return "" +} + +// CommitLease promotes a reservation to a running job and hands over the +// secrets. The executor then runs the real engine with the slot it already +// holds. +type CommitLease struct { + state protoimpl.MessageState `protogen:"open.v1"` + LeaseId string `protobuf:"bytes,1,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + Secrets []*Secret `protobuf:"bytes,2,rep,name=secrets,proto3" json:"secrets,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *CommitLease) Reset() { + *x = CommitLease{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[7] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *CommitLease) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*CommitLease) ProtoMessage() {} + +func (x *CommitLease) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[7] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use CommitLease.ProtoReflect.Descriptor instead. +func (*CommitLease) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{7} +} + +func (x *CommitLease) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +func (x *CommitLease) GetSecrets() []*Secret { + if x != nil { + return x.Secrets + } + return nil +} + +// Committed acks a CommitLease. +type Committed struct { + state protoimpl.MessageState `protogen:"open.v1"` + LeaseId string `protobuf:"bytes,1,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *Committed) Reset() { + *x = Committed{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[8] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Committed) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Committed) ProtoMessage() {} + +func (x *Committed) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[8] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Committed.ProtoReflect.Descriptor instead. +func (*Committed) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{8} +} + +func (x *Committed) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +// ReleaseLease tells the executor to drop a reservation it never committed (the +// mill picked another node, or is cleaning up). +type ReleaseLease struct { + state protoimpl.MessageState `protogen:"open.v1"` + LeaseId string `protobuf:"bytes,1,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ReleaseLease) Reset() { + *x = ReleaseLease{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[9] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ReleaseLease) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ReleaseLease) ProtoMessage() {} + +func (x *ReleaseLease) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[9] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ReleaseLease.ProtoReflect.Descriptor instead. +func (*ReleaseLease) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{9} +} + +func (x *ReleaseLease) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +// CancelAttempt cancels a running attempt (user cancel / DestroyWorkflow). +type CancelAttempt struct { + state protoimpl.MessageState `protogen:"open.v1"` + LeaseId string `protobuf:"bytes,1,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + Reason string `protobuf:"bytes,2,opt,name=reason,proto3" json:"reason,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *CancelAttempt) Reset() { + *x = CancelAttempt{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[10] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *CancelAttempt) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*CancelAttempt) ProtoMessage() {} + +func (x *CancelAttempt) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[10] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use CancelAttempt.ProtoReflect.Descriptor instead. +func (*CancelAttempt) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{10} +} + +func (x *CancelAttempt) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +func (x *CancelAttempt) GetReason() string { + if x != nil { + return x.Reason + } + return "" +} + +// StatusEvent relays a non-terminal status (i.e. running) the executor wrote to +// its own eventstream. Terminals are authored by the mill from the attempt +// result, so they are never relayed here. Carries a per-session monotonic +// offset for gap-free replay across reconnects. +type StatusEvent struct { + state protoimpl.MessageState `protogen:"open.v1"` + Offset uint64 `protobuf:"varint,1,opt,name=offset,proto3" json:"offset,omitempty"` + LeaseId string `protobuf:"bytes,2,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + Status string `protobuf:"bytes,3,opt,name=status,proto3" json:"status,omitempty"` + Error string `protobuf:"bytes,4,opt,name=error,proto3" json:"error,omitempty"` + ExitCode int64 `protobuf:"varint,5,opt,name=exit_code,json=exitCode,proto3" json:"exit_code,omitempty"` + Workflow string `protobuf:"bytes,6,opt,name=workflow,proto3" json:"workflow,omitempty"` + PipelineAturi string `protobuf:"bytes,7,opt,name=pipeline_aturi,json=pipelineAturi,proto3" json:"pipeline_aturi,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *StatusEvent) Reset() { + *x = StatusEvent{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[11] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *StatusEvent) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*StatusEvent) ProtoMessage() {} + +func (x *StatusEvent) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[11] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use StatusEvent.ProtoReflect.Descriptor instead. +func (*StatusEvent) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{11} +} + +func (x *StatusEvent) GetOffset() uint64 { + if x != nil { + return x.Offset + } + return 0 +} + +func (x *StatusEvent) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +func (x *StatusEvent) GetStatus() string { + if x != nil { + return x.Status + } + return "" +} + +func (x *StatusEvent) GetError() string { + if x != nil { + return x.Error + } + return "" +} + +func (x *StatusEvent) GetExitCode() int64 { + if x != nil { + return x.ExitCode + } + return 0 +} + +func (x *StatusEvent) GetWorkflow() string { + if x != nil { + return x.Workflow + } + return "" +} + +func (x *StatusEvent) GetPipelineAturi() string { + if x != nil { + return x.PipelineAturi + } + return "" +} + +// LogLine relays one already-encoded models.LogLine JSON line. Carries an +// offset on the same per-session sequence as StatusEvent. +type LogLine struct { + state protoimpl.MessageState `protogen:"open.v1"` + Offset uint64 `protobuf:"varint,1,opt,name=offset,proto3" json:"offset,omitempty"` + LeaseId string `protobuf:"bytes,2,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + RawJson []byte `protobuf:"bytes,3,opt,name=raw_json,json=rawJson,proto3" json:"raw_json,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *LogLine) Reset() { + *x = LogLine{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[12] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *LogLine) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*LogLine) ProtoMessage() {} + +func (x *LogLine) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[12] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use LogLine.ProtoReflect.Descriptor instead. +func (*LogLine) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{12} +} + +func (x *LogLine) GetOffset() uint64 { + if x != nil { + return x.Offset + } + return 0 +} + +func (x *LogLine) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +func (x *LogLine) GetRawJson() []byte { + if x != nil { + return x.RawJson + } + return nil +} + +// AttemptResult is the terminal signal that wakes the mill's blocked RunStep +// and resolves the lease. terminal_status is one of success/failed/timeout/ +// cancelled. Carries an offset on the same per-session sequence. +type AttemptResult struct { + state protoimpl.MessageState `protogen:"open.v1"` + Offset uint64 `protobuf:"varint,1,opt,name=offset,proto3" json:"offset,omitempty"` + LeaseId string `protobuf:"bytes,2,opt,name=lease_id,json=leaseId,proto3" json:"lease_id,omitempty"` + TerminalStatus string `protobuf:"bytes,3,opt,name=terminal_status,json=terminalStatus,proto3" json:"terminal_status,omitempty"` + Error string `protobuf:"bytes,4,opt,name=error,proto3" json:"error,omitempty"` + ExitCode int64 `protobuf:"varint,5,opt,name=exit_code,json=exitCode,proto3" json:"exit_code,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *AttemptResult) Reset() { + *x = AttemptResult{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[13] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *AttemptResult) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*AttemptResult) ProtoMessage() {} + +func (x *AttemptResult) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[13] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use AttemptResult.ProtoReflect.Descriptor instead. +func (*AttemptResult) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{13} +} + +func (x *AttemptResult) GetOffset() uint64 { + if x != nil { + return x.Offset + } + return 0 +} + +func (x *AttemptResult) GetLeaseId() string { + if x != nil { + return x.LeaseId + } + return "" +} + +func (x *AttemptResult) GetTerminalStatus() string { + if x != nil { + return x.TerminalStatus + } + return "" +} + +func (x *AttemptResult) GetError() string { + if x != nil { + return x.Error + } + return "" +} + +func (x *AttemptResult) GetExitCode() int64 { + if x != nil { + return x.ExitCode + } + return 0 +} + +// Ack tells the executor the mill has durably processed all relay entries up +// to and including up_to_offset, so it may trim its buffer. +type Ack struct { + state protoimpl.MessageState `protogen:"open.v1"` + UpToOffset uint64 `protobuf:"varint,1,opt,name=up_to_offset,json=upToOffset,proto3" json:"up_to_offset,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *Ack) Reset() { + *x = Ack{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[14] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Ack) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Ack) ProtoMessage() {} + +func (x *Ack) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[14] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Ack.ProtoReflect.Descriptor instead. +func (*Ack) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{14} +} + +func (x *Ack) GetUpToOffset() uint64 { + if x != nil { + return x.UpToOffset + } + return 0 +} + +type Message struct { + state protoimpl.MessageState `protogen:"open.v1"` + Hello *Hello `protobuf:"bytes,1,opt,name=hello,proto3" json:"hello,omitempty"` + Resume *Resume `protobuf:"bytes,2,opt,name=resume,proto3" json:"resume,omitempty"` + NodeSnapshot *NodeSnapshot `protobuf:"bytes,3,opt,name=node_snapshot,json=nodeSnapshot,proto3" json:"node_snapshot,omitempty"` + ReserveSeat *ReserveSeat `protobuf:"bytes,4,opt,name=reserve_seat,json=reserveSeat,proto3" json:"reserve_seat,omitempty"` + ReserveResult *ReserveResult `protobuf:"bytes,5,opt,name=reserve_result,json=reserveResult,proto3" json:"reserve_result,omitempty"` + CommitLease *CommitLease `protobuf:"bytes,6,opt,name=commit_lease,json=commitLease,proto3" json:"commit_lease,omitempty"` + Committed *Committed `protobuf:"bytes,7,opt,name=committed,proto3" json:"committed,omitempty"` + ReleaseLease *ReleaseLease `protobuf:"bytes,8,opt,name=release_lease,json=releaseLease,proto3" json:"release_lease,omitempty"` + CancelAttempt *CancelAttempt `protobuf:"bytes,9,opt,name=cancel_attempt,json=cancelAttempt,proto3" json:"cancel_attempt,omitempty"` + StatusEvent *StatusEvent `protobuf:"bytes,10,opt,name=status_event,json=statusEvent,proto3" json:"status_event,omitempty"` + LogLine *LogLine `protobuf:"bytes,11,opt,name=log_line,json=logLine,proto3" json:"log_line,omitempty"` + AttemptResult *AttemptResult `protobuf:"bytes,12,opt,name=attempt_result,json=attemptResult,proto3" json:"attempt_result,omitempty"` + Ack *Ack `protobuf:"bytes,13,opt,name=ack,proto3" json:"ack,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *Message) Reset() { + *x = Message{} + mi := &file_spindle_mill_v1_mill_proto_msgTypes[15] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Message) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Message) ProtoMessage() {} + +func (x *Message) ProtoReflect() protoreflect.Message { + mi := &file_spindle_mill_v1_mill_proto_msgTypes[15] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Message.ProtoReflect.Descriptor instead. +func (*Message) Descriptor() ([]byte, []int) { + return file_spindle_mill_v1_mill_proto_rawDescGZIP(), []int{15} +} + +func (x *Message) GetHello() *Hello { + if x != nil { + return x.Hello + } + return nil +} + +func (x *Message) GetResume() *Resume { + if x != nil { + return x.Resume + } + return nil +} + +func (x *Message) GetNodeSnapshot() *NodeSnapshot { + if x != nil { + return x.NodeSnapshot + } + return nil +} + +func (x *Message) GetReserveSeat() *ReserveSeat { + if x != nil { + return x.ReserveSeat + } + return nil +} + +func (x *Message) GetReserveResult() *ReserveResult { + if x != nil { + return x.ReserveResult + } + return nil +} + +func (x *Message) GetCommitLease() *CommitLease { + if x != nil { + return x.CommitLease + } + return nil +} + +func (x *Message) GetCommitted() *Committed { + if x != nil { + return x.Committed + } + return nil +} + +func (x *Message) GetReleaseLease() *ReleaseLease { + if x != nil { + return x.ReleaseLease + } + return nil +} + +func (x *Message) GetCancelAttempt() *CancelAttempt { + if x != nil { + return x.CancelAttempt + } + return nil +} + +func (x *Message) GetStatusEvent() *StatusEvent { + if x != nil { + return x.StatusEvent + } + return nil +} + +func (x *Message) GetLogLine() *LogLine { + if x != nil { + return x.LogLine + } + return nil +} + +func (x *Message) GetAttemptResult() *AttemptResult { + if x != nil { + return x.AttemptResult + } + return nil +} + +func (x *Message) GetAck() *Ack { + if x != nil { + return x.Ack + } + return nil +} + +var File_spindle_mill_v1_mill_proto protoreflect.FileDescriptor + +const file_spindle_mill_v1_mill_proto_rawDesc = "" + + "\n" + + "\x1aspindle/mill/v1/mill.proto\x12\x0fspindle.mill.v1\x1a\x1bbuf/validate/validate.proto\"\xa3\x01\n" + + "\x05Hello\x12)\n" + + "\x10protocol_version\x18\x01 \x01(\rR\x0fprotocolVersion\x12 \n" + + "\anode_id\x18\x02 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\x06nodeId\x12\x18\n" + + "\aengines\x18\x03 \x03(\tR\aengines\x12\x12\n" + + "\x04arch\x18\x04 \x01(\tR\x04arch\x12\x1f\n" + + "\vlast_offset\x18\x05 \x01(\x04R\n" + + "lastOffset\"'\n" + + "\x06Resume\x12\x1d\n" + + "\n" + + "ack_offset\x18\x01 \x01(\x04R\tackOffset\"\x9a\x01\n" + + "\x0eEngineSnapshot\x12\x1d\n" + + "\n" + + "free_seats\x18\x01 \x01(\rR\tfreeSeats\x12&\n" + + "\x0ffree_memory_mib\x18\x02 \x01(\x03R\rfreeMemoryMib\x12\x1d\n" + + "\n" + + "free_vcpus\x18\x03 \x01(\x03R\tfreeVcpus\x12\"\n" + + "\rfree_disk_mib\x18\x04 \x01(\x03R\vfreeDiskMib\"\xdc\x01\n" + + "\fNodeSnapshot\x12\x17\n" + + "\anode_id\x18\x01 \x01(\tR\x06nodeId\x12\x10\n" + + "\x03seq\x18\x02 \x01(\x04R\x03seq\x12D\n" + + "\aengines\x18\x03 \x03(\v2*.spindle.mill.v1.NodeSnapshot.EnginesEntryR\aengines\x1a[\n" + + "\fEnginesEntry\x12\x10\n" + + "\x03key\x18\x01 \x01(\tR\x03key\x125\n" + + "\x05value\x18\x02 \x01(\v2\x1f.spindle.mill.v1.EngineSnapshotR\x05value:\x028\x01\"\x80\x02\n" + + "\vReserveSeat\x12\"\n" + + "\blease_id\x18\x01 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x12,\n" + + "\rtarget_engine\x18\x02 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\ftargetEngine\x12*\n" + + "\x11raw_pipeline_json\x18\x03 \x01(\tR\x0frawPipelineJson\x12*\n" + + "\x11raw_workflow_json\x18\x04 \x01(\tR\x0frawWorkflowJson\x12\x12\n" + + "\x04knot\x18\x05 \x01(\tR\x04knot\x12\x12\n" + + "\x04rkey\x18\x06 \x01(\tR\x04rkey\x12\x1f\n" + + "\vttl_seconds\x18\a \x01(\rR\n" + + "ttlSeconds\"\x8a\x01\n" + + "\rReserveResult\x12\"\n" + + "\blease_id\x18\x01 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x12\x1a\n" + + "\baccepted\x18\x02 \x01(\bR\baccepted\x12#\n" + + "\rreject_reason\x18\x03 \x01(\tR\frejectReason\x12\x14\n" + + "\x05score\x18\x04 \x01(\x01R\x05score\"0\n" + + "\x06Secret\x12\x10\n" + + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + + "\x05value\x18\x02 \x01(\tR\x05value\"d\n" + + "\vCommitLease\x12\"\n" + + "\blease_id\x18\x01 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x121\n" + + "\asecrets\x18\x02 \x03(\v2\x17.spindle.mill.v1.SecretR\asecrets\"/\n" + + "\tCommitted\x12\"\n" + + "\blease_id\x18\x01 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\"2\n" + + "\fReleaseLease\x12\"\n" + + "\blease_id\x18\x01 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\"K\n" + + "\rCancelAttempt\x12\"\n" + + "\blease_id\x18\x01 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x12\x16\n" + + "\x06reason\x18\x02 \x01(\tR\x06reason\"\xd7\x01\n" + + "\vStatusEvent\x12\x16\n" + + "\x06offset\x18\x01 \x01(\x04R\x06offset\x12\"\n" + + "\blease_id\x18\x02 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x12\x16\n" + + "\x06status\x18\x03 \x01(\tR\x06status\x12\x14\n" + + "\x05error\x18\x04 \x01(\tR\x05error\x12\x1b\n" + + "\texit_code\x18\x05 \x01(\x03R\bexitCode\x12\x1a\n" + + "\bworkflow\x18\x06 \x01(\tR\bworkflow\x12%\n" + + "\x0epipeline_aturi\x18\a \x01(\tR\rpipelineAturi\"`\n" + + "\aLogLine\x12\x16\n" + + "\x06offset\x18\x01 \x01(\x04R\x06offset\x12\"\n" + + "\blease_id\x18\x02 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x12\x19\n" + + "\braw_json\x18\x03 \x01(\fR\arawJson\"\xa7\x01\n" + + "\rAttemptResult\x12\x16\n" + + "\x06offset\x18\x01 \x01(\x04R\x06offset\x12\"\n" + + "\blease_id\x18\x02 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x12'\n" + + "\x0fterminal_status\x18\x03 \x01(\tR\x0eterminalStatus\x12\x14\n" + + "\x05error\x18\x04 \x01(\tR\x05error\x12\x1b\n" + + "\texit_code\x18\x05 \x01(\x03R\bexitCode\"'\n" + + "\x03Ack\x12 \n" + + "\fup_to_offset\x18\x01 \x01(\x04R\n" + + "upToOffset\"\xcc\a\n" + + "\aMessage\x12,\n" + + "\x05hello\x18\x01 \x01(\v2\x16.spindle.mill.v1.HelloR\x05hello\x12/\n" + + "\x06resume\x18\x02 \x01(\v2\x17.spindle.mill.v1.ResumeR\x06resume\x12B\n" + + "\rnode_snapshot\x18\x03 \x01(\v2\x1d.spindle.mill.v1.NodeSnapshotR\fnodeSnapshot\x12?\n" + + "\freserve_seat\x18\x04 \x01(\v2\x1c.spindle.mill.v1.ReserveSeatR\vreserveSeat\x12E\n" + + "\x0ereserve_result\x18\x05 \x01(\v2\x1e.spindle.mill.v1.ReserveResultR\rreserveResult\x12?\n" + + "\fcommit_lease\x18\x06 \x01(\v2\x1c.spindle.mill.v1.CommitLeaseR\vcommitLease\x128\n" + + "\tcommitted\x18\a \x01(\v2\x1a.spindle.mill.v1.CommittedR\tcommitted\x12B\n" + + "\rrelease_lease\x18\b \x01(\v2\x1d.spindle.mill.v1.ReleaseLeaseR\freleaseLease\x12E\n" + + "\x0ecancel_attempt\x18\t \x01(\v2\x1e.spindle.mill.v1.CancelAttemptR\rcancelAttempt\x12?\n" + + "\fstatus_event\x18\n" + + " \x01(\v2\x1c.spindle.mill.v1.StatusEventR\vstatusEvent\x123\n" + + "\blog_line\x18\v \x01(\v2\x18.spindle.mill.v1.LogLineR\alogLine\x12E\n" + + "\x0eattempt_result\x18\f \x01(\v2\x1e.spindle.mill.v1.AttemptResultR\rattemptResult\x12&\n" + + "\x03ack\x18\r \x01(\v2\x14.spindle.mill.v1.AckR\x03ack:\xaa\x01\xbaH\xa6\x01\"\xa3\x01\n" + + "\x05hello\n" + + "\x06resume\n" + + "\rnode_snapshot\n" + + "\freserve_seat\n" + + "\x0ereserve_result\n" + + "\fcommit_lease\n" + + "\tcommitted\n" + + "\rrelease_lease\n" + + "\x0ecancel_attempt\n" + + "\fstatus_event\n" + + "\blog_line\n" + + "\x0eattempt_result\n" + + "\x03ack\x10\x01B0Z.tangled.org/core/spindle/mill/proto/gen;millv1b\x06proto3" + +var ( + file_spindle_mill_v1_mill_proto_rawDescOnce sync.Once + file_spindle_mill_v1_mill_proto_rawDescData []byte +) + +func file_spindle_mill_v1_mill_proto_rawDescGZIP() []byte { + file_spindle_mill_v1_mill_proto_rawDescOnce.Do(func() { + file_spindle_mill_v1_mill_proto_rawDescData = protoimpl.X.CompressGZIP(unsafe.Slice(unsafe.StringData(file_spindle_mill_v1_mill_proto_rawDesc), len(file_spindle_mill_v1_mill_proto_rawDesc))) + }) + return file_spindle_mill_v1_mill_proto_rawDescData +} + +var file_spindle_mill_v1_mill_proto_msgTypes = make([]protoimpl.MessageInfo, 17) +var file_spindle_mill_v1_mill_proto_goTypes = []any{ + (*Hello)(nil), // 0: spindle.mill.v1.Hello + (*Resume)(nil), // 1: spindle.mill.v1.Resume + (*EngineSnapshot)(nil), // 2: spindle.mill.v1.EngineSnapshot + (*NodeSnapshot)(nil), // 3: spindle.mill.v1.NodeSnapshot + (*ReserveSeat)(nil), // 4: spindle.mill.v1.ReserveSeat + (*ReserveResult)(nil), // 5: spindle.mill.v1.ReserveResult + (*Secret)(nil), // 6: spindle.mill.v1.Secret + (*CommitLease)(nil), // 7: spindle.mill.v1.CommitLease + (*Committed)(nil), // 8: spindle.mill.v1.Committed + (*ReleaseLease)(nil), // 9: spindle.mill.v1.ReleaseLease + (*CancelAttempt)(nil), // 10: spindle.mill.v1.CancelAttempt + (*StatusEvent)(nil), // 11: spindle.mill.v1.StatusEvent + (*LogLine)(nil), // 12: spindle.mill.v1.LogLine + (*AttemptResult)(nil), // 13: spindle.mill.v1.AttemptResult + (*Ack)(nil), // 14: spindle.mill.v1.Ack + (*Message)(nil), // 15: spindle.mill.v1.Message + nil, // 16: spindle.mill.v1.NodeSnapshot.EnginesEntry +} +var file_spindle_mill_v1_mill_proto_depIdxs = []int32{ + 16, // 0: spindle.mill.v1.NodeSnapshot.engines:type_name -> spindle.mill.v1.NodeSnapshot.EnginesEntry + 6, // 1: spindle.mill.v1.CommitLease.secrets:type_name -> spindle.mill.v1.Secret + 0, // 2: spindle.mill.v1.Message.hello:type_name -> spindle.mill.v1.Hello + 1, // 3: spindle.mill.v1.Message.resume:type_name -> spindle.mill.v1.Resume + 3, // 4: spindle.mill.v1.Message.node_snapshot:type_name -> spindle.mill.v1.NodeSnapshot + 4, // 5: spindle.mill.v1.Message.reserve_seat:type_name -> spindle.mill.v1.ReserveSeat + 5, // 6: spindle.mill.v1.Message.reserve_result:type_name -> spindle.mill.v1.ReserveResult + 7, // 7: spindle.mill.v1.Message.commit_lease:type_name -> spindle.mill.v1.CommitLease + 8, // 8: spindle.mill.v1.Message.committed:type_name -> spindle.mill.v1.Committed + 9, // 9: spindle.mill.v1.Message.release_lease:type_name -> spindle.mill.v1.ReleaseLease + 10, // 10: spindle.mill.v1.Message.cancel_attempt:type_name -> spindle.mill.v1.CancelAttempt + 11, // 11: spindle.mill.v1.Message.status_event:type_name -> spindle.mill.v1.StatusEvent + 12, // 12: spindle.mill.v1.Message.log_line:type_name -> spindle.mill.v1.LogLine + 13, // 13: spindle.mill.v1.Message.attempt_result:type_name -> spindle.mill.v1.AttemptResult + 14, // 14: spindle.mill.v1.Message.ack:type_name -> spindle.mill.v1.Ack + 2, // 15: spindle.mill.v1.NodeSnapshot.EnginesEntry.value:type_name -> spindle.mill.v1.EngineSnapshot + 16, // [16:16] is the sub-list for method output_type + 16, // [16:16] is the sub-list for method input_type + 16, // [16:16] is the sub-list for extension type_name + 16, // [16:16] is the sub-list for extension extendee + 0, // [0:16] is the sub-list for field type_name +} + +func init() { file_spindle_mill_v1_mill_proto_init() } +func file_spindle_mill_v1_mill_proto_init() { + if File_spindle_mill_v1_mill_proto != nil { + return + } + type x struct{} + out := protoimpl.TypeBuilder{ + File: protoimpl.DescBuilder{ + GoPackagePath: reflect.TypeOf(x{}).PkgPath(), + RawDescriptor: unsafe.Slice(unsafe.StringData(file_spindle_mill_v1_mill_proto_rawDesc), len(file_spindle_mill_v1_mill_proto_rawDesc)), + NumEnums: 0, + NumMessages: 17, + NumExtensions: 0, + NumServices: 0, + }, + GoTypes: file_spindle_mill_v1_mill_proto_goTypes, + DependencyIndexes: file_spindle_mill_v1_mill_proto_depIdxs, + MessageInfos: file_spindle_mill_v1_mill_proto_msgTypes, + }.Build() + File_spindle_mill_v1_mill_proto = out.File + file_spindle_mill_v1_mill_proto_goTypes = nil + file_spindle_mill_v1_mill_proto_depIdxs = nil +} diff --git a/spindle/mill/proto/protocol.go b/spindle/mill/proto/protocol.go new file mode 100644 index 00000000..7cdda436 --- /dev/null +++ b/spindle/mill/proto/protocol.go @@ -0,0 +1,103 @@ +// Package millproto carries the mill<->executor session protocol: a new +// message vocabulary over the same length-prefixed protobuf framing the spindle +// already uses to talk to the microVM guest (see spindle/agentproto). Only the +// framing pattern is shared; the messages are entirely separate. +package millproto + +import ( + "encoding/binary" + "fmt" + "io" + "sync" + + "buf.build/go/protovalidate" + "google.golang.org/protobuf/proto" + + millv1 "tangled.org/core/spindle/mill/proto/gen" +) + +const ( + ProtocolVersion = 1 + // generous vs agentproto's 1 MiB: a ReserveSeat carries the raw pipeline + + // workflow JSON, and relayed log lines can be chunky. + MaxMessageBytes = 8 * 1024 * 1024 +) + +type Message = millv1.Message + +var validator protovalidate.Validator + +func init() { + var err error + validator, err = protovalidate.New() + if err != nil { + panic(fmt.Errorf("failed to initialize protovalidate validator: %w", err)) + } +} + +type Encoder struct { + mu sync.Mutex + w io.Writer +} + +func NewEncoder(w io.Writer) *Encoder { + return &Encoder{w: w} +} + +func (e *Encoder) Encode(msg *Message) error { + if err := validator.Validate(msg); err != nil { + return fmt.Errorf("validate fleet message: %w", err) + } + + data, err := proto.Marshal(msg) + if err != nil { + return fmt.Errorf("marshal fleet message: %w", err) + } + if len(data) > MaxMessageBytes { + return fmt.Errorf("fleet message exceeded %d bytes", MaxMessageBytes) + } + + // single write of header+payload so this maps to exactly one websocket + // binary frame when the writer is a ws stream. + frame := make([]byte, 4+len(data)) + binary.BigEndian.PutUint32(frame[:4], uint32(len(data))) + copy(frame[4:], data) + + e.mu.Lock() + defer e.mu.Unlock() + _, err = e.w.Write(frame) + return err +} + +type Decoder struct { + r io.Reader +} + +func NewDecoder(r io.Reader) *Decoder { + return &Decoder{r: r} +} + +func (d *Decoder) Decode() (*Message, error) { + msg := &Message{} + var header [4]byte + if _, err := io.ReadFull(d.r, header[:]); err != nil { + return msg, err + } + + size := binary.BigEndian.Uint32(header[:]) + if size > MaxMessageBytes { + return msg, fmt.Errorf("fleet message exceeded %d bytes", MaxMessageBytes) + } + + data := make([]byte, size) + if _, err := io.ReadFull(d.r, data); err != nil { + return msg, err + } + if err := proto.Unmarshal(data, msg); err != nil { + return msg, fmt.Errorf("parse fleet message: %w", err) + } + if err := validator.Validate(msg); err != nil { + return msg, fmt.Errorf("validate fleet message: %w", err) + } + return msg, nil +} diff --git a/spindle/mill/proto/protocol_test.go b/spindle/mill/proto/protocol_test.go new file mode 100644 index 00000000..a3db7d23 --- /dev/null +++ b/spindle/mill/proto/protocol_test.go @@ -0,0 +1,66 @@ +package millproto + +import ( + "bytes" + "encoding/binary" + "testing" + + millv1 "tangled.org/core/spindle/mill/proto/gen" +) + +func TestEncodeDecodeRoundTrip(t *testing.T) { + var buf bytes.Buffer + enc := NewEncoder(&buf) + + want := &Message{ + ReserveSeat: &millv1.ReserveSeat{ + LeaseId: "lease-1", + TargetEngine: "microvm", + RawWorkflowJson: `{"name":"build"}`, + Knot: "knot.example", + Rkey: "abc123", + TtlSeconds: 30, + }, + } + if err := enc.Encode(want); err != nil { + t.Fatalf("Encode() error = %v", err) + } + + got, err := NewDecoder(&buf).Decode() + if err != nil { + t.Fatalf("Decode() error = %v", err) + } + rs := got.GetReserveSeat() + if rs == nil { + t.Fatal("decoded message missing reserve_seat") + } + if rs.LeaseId != "lease-1" || rs.TargetEngine != "microvm" || rs.TtlSeconds != 30 { + t.Fatalf("round-trip mismatch: %+v", rs) + } +} + +func TestDecoderRejectsOversizedMessage(t *testing.T) { + var tooLarge bytes.Buffer + var header [4]byte + binary.BigEndian.PutUint32(header[:], MaxMessageBytes+1) + tooLarge.Write(header[:]) + + if _, err := NewDecoder(&tooLarge).Decode(); err == nil { + t.Fatal("expected oversized message error") + } +} + +func TestValidationOneofRequired(t *testing.T) { + if err := validator.Validate(&Message{Ack: &millv1.Ack{UpToOffset: 5}}); err != nil { + t.Fatalf("expected valid message to pass, got: %v", err) + } + if err := validator.Validate(&Message{}); err == nil { + t.Fatal("expected message with zero payloads to fail validation") + } + if err := validator.Validate(&Message{ + Ack: &millv1.Ack{}, + Committed: &millv1.Committed{LeaseId: "x"}, + }); err == nil { + t.Fatal("expected message with multiple payloads to fail validation") + } +} diff --git a/spindle/mill/proto/spindle/mill/v1/mill.proto b/spindle/mill/proto/spindle/mill/v1/mill.proto new file mode 100644 index 00000000..e53bb793 --- /dev/null +++ b/spindle/mill/proto/spindle/mill/v1/mill.proto @@ -0,0 +1,169 @@ +syntax = "proto3"; + +package spindle.mill.v1; + +import "buf/validate/validate.proto"; + +option go_package = "tangled.org/core/spindle/mill/proto/gen;millv1"; + +// Hello is the first frame an executor sends after dialing the mill. It +// carries the static traits of the node plus a resume hint. +message Hello { + uint32 protocol_version = 1; + // stable across reconnects (e.g. hostname / DID); the mill keys sessions on + // this so a brief blip reattaches the same node rather than creating a new one. + string node_id = 2 [(buf.validate.field).string.min_len = 1]; + // engine names this node can run ("microvm", "nixery"). + repeated string engines = 3; + // GOARCH of the node, so the mill won't place arch-incompatible jobs. + string arch = 4; + // the highest relay offset the executor believes it has sent; a resume hint. + uint64 last_offset = 5; +} + +// Resume is the mill's reply to Hello. The executor replays every buffered +// relay entry with offset strictly greater than ack_offset before sending new +// ones. +message Resume { + uint64 ack_offset = 1; +} + +// EngineSnapshot is the changing per-engine state on a node. free_seats is a +// coarse "can you take more" hint; 0 also means draining (the mill treats a +// draining node as a full one until it leaves). The resource fields are coarse +// budget headroom for ranking only; the executor makes the real yes/no call in +// ReserveResult. +message EngineSnapshot { + uint32 free_seats = 1; + int64 free_memory_mib = 2; + int64 free_vcpus = 3; + int64 free_disk_mib = 4; +} + +// NodeSnapshot is pushed on connect, periodically, and right after any state +// change (reserve, commit, terminal). +message NodeSnapshot { + string node_id = 1; + uint64 seq = 2; + map engines = 3; +} + +// ReserveSeat asks an executor to hold a seat for a job. Zero secrets ride this +// message; the raw pipeline/workflow are carried as JSON since processPipeline +// already round-trips them through JSON. +message ReserveSeat { + string lease_id = 1 [(buf.validate.field).string.min_len = 1]; + string target_engine = 2 [(buf.validate.field).string.min_len = 1]; + string raw_pipeline_json = 3; + string raw_workflow_json = 4; + // the pipeline id (knot + rkey); the executor reconstructs the exact + // WorkflowId so its relayed status rows and log path match what the mill + // authored for "pending". + string knot = 5; + string rkey = 6; + uint32 ttl_seconds = 7; +} + +// ReserveResult is the executor's accept/reject for a ReserveSeat. +message ReserveResult { + string lease_id = 1 [(buf.validate.field).string.min_len = 1]; + bool accepted = 2; + string reject_reason = 3; + // optional bid score the mill ranks accepted leases by (higher is better). + double score = 4; +} + +// Secret is a single unlocked secret, sent only inside CommitLease. +message Secret { + string key = 1; + string value = 2; +} + +// CommitLease promotes a reservation to a running job and hands over the +// secrets. The executor then runs the real engine with the slot it already +// holds. +message CommitLease { + string lease_id = 1 [(buf.validate.field).string.min_len = 1]; + repeated Secret secrets = 2; +} + +// Committed acks a CommitLease. +message Committed { + string lease_id = 1 [(buf.validate.field).string.min_len = 1]; +} + +// ReleaseLease tells the executor to drop a reservation it never committed (the +// mill picked another node, or is cleaning up). +message ReleaseLease { + string lease_id = 1 [(buf.validate.field).string.min_len = 1]; +} + +// CancelAttempt cancels a running attempt (user cancel / DestroyWorkflow). +message CancelAttempt { + string lease_id = 1 [(buf.validate.field).string.min_len = 1]; + string reason = 2; +} + +// StatusEvent relays a non-terminal status (i.e. running) the executor wrote to +// its own eventstream. Terminals are authored by the mill from the attempt +// result, so they are never relayed here. Carries a per-session monotonic +// offset for gap-free replay across reconnects. +message StatusEvent { + uint64 offset = 1; + string lease_id = 2 [(buf.validate.field).string.min_len = 1]; + string status = 3; + string error = 4; + int64 exit_code = 5; + string workflow = 6; + string pipeline_aturi = 7; +} + +// LogLine relays one already-encoded models.LogLine JSON line. Carries an +// offset on the same per-session sequence as StatusEvent. +message LogLine { + uint64 offset = 1; + string lease_id = 2 [(buf.validate.field).string.min_len = 1]; + bytes raw_json = 3; +} + +// AttemptResult is the terminal signal that wakes the mill's blocked RunStep +// and resolves the lease. terminal_status is one of success/failed/timeout/ +// cancelled. Carries an offset on the same per-session sequence. +message AttemptResult { + uint64 offset = 1; + string lease_id = 2 [(buf.validate.field).string.min_len = 1]; + string terminal_status = 3; + string error = 4; + int64 exit_code = 5; +} + +// Ack tells the executor the mill has durably processed all relay entries up +// to and including up_to_offset, so it may trim its buffer. +message Ack { + uint64 up_to_offset = 1; +} + +message Message { + option (buf.validate.message).oneof = { + fields: [ + "hello", "resume", "node_snapshot", "reserve_seat", "reserve_result", + "commit_lease", "committed", "release_lease", "cancel_attempt", + "status_event", "log_line", "attempt_result", "ack" + ], + required: true + }; + + Hello hello = 1; + Resume resume = 2; + NodeSnapshot node_snapshot = 3; + ReserveSeat reserve_seat = 4; + ReserveResult reserve_result = 5; + CommitLease commit_lease = 6; + Committed committed = 7; + ReleaseLease release_lease = 8; + CancelAttempt cancel_attempt = 9; + StatusEvent status_event = 10; + LogLine log_line = 11; + AttemptResult attempt_result = 12; + Ack ack = 13; +} diff --git a/spindle/mill/proto/ws.go b/spindle/mill/proto/ws.go new file mode 100644 index 00000000..6fdc0c7c --- /dev/null +++ b/spindle/mill/proto/ws.go @@ -0,0 +1,58 @@ +package millproto + +import ( + "io" + "sync" + + "github.com/gorilla/websocket" +) + +// WSStream adapts a gorilla websocket connection to an io.ReadWriteCloser so the +// length-prefixed fleet framing rides over it. Each Encode produces exactly one +// binary frame; the reader reassembles the byte stream across frames. +type WSStream struct { + conn *websocket.Conn + + rmu sync.Mutex + r io.Reader // current message reader, advanced as frames are consumed + + wmu sync.Mutex +} + +func NewWSStream(conn *websocket.Conn) *WSStream { + return &WSStream{conn: conn} +} + +func (s *WSStream) Read(p []byte) (int, error) { + s.rmu.Lock() + defer s.rmu.Unlock() + for { + if s.r == nil { + _, r, err := s.conn.NextReader() + if err != nil { + return 0, err + } + s.r = r + } + n, err := s.r.Read(p) + if err == io.EOF { + s.r = nil + if n > 0 { + return n, nil + } + continue + } + return n, err + } +} + +func (s *WSStream) Write(p []byte) (int, error) { + s.wmu.Lock() + defer s.wmu.Unlock() + if err := s.conn.WriteMessage(websocket.BinaryMessage, p); err != nil { + return 0, err + } + return len(p), nil +} + +func (s *WSStream) Close() error { return s.conn.Close() } -- 2.51.2