From 4248f28cb6c319d6cd66ab4287939935dd36b90d Mon Sep 17 00:00:00 2001 From: Dhamith Hewamullage Date: Tue, 29 Sep 2026 19:16:22 +0530 Subject: [PATCH 01/10] add rpcs for forecasts, process usage and snapshots --- internal/api/api.pb.go | 550 ++++++++++++++++++++++++++++++++---- internal/api/api.proto | 50 ++++ internal/api/api_grpc.pb.go | 116 ++++++++ 3 files changed, 662 insertions(+), 54 deletions(-) diff --git a/internal/api/api.pb.go b/internal/api/api.pb.go index 2f3abc1..7ff53b4 100644 --- a/internal/api/api.pb.go +++ b/internal/api/api.pb.go @@ -328,7 +328,9 @@ type HostSummary struct { // highest severity of the open alerts, 0 when there are none WorstSeverity int32 `protobuf:"varint,14,opt,name=worstSeverity,proto3" json:"worstSeverity,omitempty"` // running containers at the latest snapshot - Containers int32 `protobuf:"varint,15,opt,name=containers,proto3" json:"containers,omitempty"` + Containers int32 `protobuf:"varint,15,opt,name=containers,proto3" json:"containers,omitempty"` + // days until the first disk fills up, unset when none is filling up + DiskFullDays *float64 `protobuf:"fixed64,16,opt,name=diskFullDays,proto3,oneof" json:"diskFullDays,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -468,6 +470,13 @@ func (x *HostSummary) GetContainers() int32 { return 0 } +func (x *HostSummary) GetDiskFullDays() float64 { + if x != nil && x.DiskFullDays != nil { + return *x.DiskFullDays + } + return 0 +} + type FleetSummary struct { state protoimpl.MessageState `protogen:"open.v1"` Hosts []*HostSummary `protobuf:"bytes,1,rep,name=hosts,proto3" json:"hosts,omitempty"` @@ -1281,6 +1290,388 @@ func (x *AlertList) GetAlerts() []*AlertRecord { return nil } +// A disk's growth over the last week and when it fills up at that rate +type DiskForecast struct { + state protoimpl.MessageState `protogen:"open.v1"` + Device string `protobuf:"bytes,1,opt,name=device,proto3" json:"device,omitempty"` + Mount string `protobuf:"bytes,2,opt,name=mount,proto3" json:"mount,omitempty"` + UsedPct float64 `protobuf:"fixed64,3,opt,name=usedPct,proto3" json:"usedPct,omitempty"` + PctPerDay float64 `protobuf:"fixed64,4,opt,name=pctPerDay,proto3" json:"pctPerDay,omitempty"` + BytesPerDay float64 `protobuf:"fixed64,5,opt,name=bytesPerDay,proto3" json:"bytesPerDay,omitempty"` + // unset when the disk is not filling up + DaysToFull *float64 `protobuf:"fixed64,6,opt,name=daysToFull,proto3,oneof" json:"daysToFull,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *DiskForecast) Reset() { + *x = DiskForecast{} + mi := &file_api_api_proto_msgTypes[20] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *DiskForecast) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*DiskForecast) ProtoMessage() {} + +func (x *DiskForecast) ProtoReflect() protoreflect.Message { + mi := &file_api_api_proto_msgTypes[20] + 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 DiskForecast.ProtoReflect.Descriptor instead. +func (*DiskForecast) Descriptor() ([]byte, []int) { + return file_api_api_proto_rawDescGZIP(), []int{20} +} + +func (x *DiskForecast) GetDevice() string { + if x != nil { + return x.Device + } + return "" +} + +func (x *DiskForecast) GetMount() string { + if x != nil { + return x.Mount + } + return "" +} + +func (x *DiskForecast) GetUsedPct() float64 { + if x != nil { + return x.UsedPct + } + return 0 +} + +func (x *DiskForecast) GetPctPerDay() float64 { + if x != nil { + return x.PctPerDay + } + return 0 +} + +func (x *DiskForecast) GetBytesPerDay() float64 { + if x != nil { + return x.BytesPerDay + } + return 0 +} + +func (x *DiskForecast) GetDaysToFull() float64 { + if x != nil && x.DaysToFull != nil { + return *x.DaysToFull + } + return 0 +} + +type DiskForecastList struct { + state protoimpl.MessageState `protogen:"open.v1"` + Disks []*DiskForecast `protobuf:"bytes,1,rep,name=disks,proto3" json:"disks,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *DiskForecastList) Reset() { + *x = DiskForecastList{} + mi := &file_api_api_proto_msgTypes[21] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *DiskForecastList) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*DiskForecastList) ProtoMessage() {} + +func (x *DiskForecastList) ProtoReflect() protoreflect.Message { + mi := &file_api_api_proto_msgTypes[21] + 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 DiskForecastList.ProtoReflect.Descriptor instead. +func (*DiskForecastList) Descriptor() ([]byte, []int) { + return file_api_api_proto_rawDescGZIP(), []int{21} +} + +func (x *DiskForecastList) GetDisks() []*DiskForecast { + if x != nil { + return x.Disks + } + return nil +} + +type ProcessUsageRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Host string `protobuf:"bytes,1,opt,name=host,proto3" json:"host,omitempty"` + From int64 `protobuf:"varint,2,opt,name=from,proto3" json:"from,omitempty"` + To int64 `protobuf:"varint,3,opt,name=to,proto3" json:"to,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ProcessUsageRequest) Reset() { + *x = ProcessUsageRequest{} + mi := &file_api_api_proto_msgTypes[22] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ProcessUsageRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ProcessUsageRequest) ProtoMessage() {} + +func (x *ProcessUsageRequest) ProtoReflect() protoreflect.Message { + mi := &file_api_api_proto_msgTypes[22] + 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 ProcessUsageRequest.ProtoReflect.Descriptor instead. +func (*ProcessUsageRequest) Descriptor() ([]byte, []int) { + return file_api_api_proto_rawDescGZIP(), []int{22} +} + +func (x *ProcessUsageRequest) GetHost() string { + if x != nil { + return x.Host + } + return "" +} + +func (x *ProcessUsageRequest) GetFrom() int64 { + if x != nil { + return x.From + } + return 0 +} + +func (x *ProcessUsageRequest) GetTo() int64 { + if x != nil { + return x.To + } + return 0 +} + +// One program's share of the range. Processes with the same name are +// added up per snapshot. +type ProcessUsage struct { + state protoimpl.MessageState `protogen:"open.v1"` + Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"` + CpuAvg float64 `protobuf:"fixed64,2,opt,name=cpuAvg,proto3" json:"cpuAvg,omitempty"` + CpuPeak float64 `protobuf:"fixed64,3,opt,name=cpuPeak,proto3" json:"cpuPeak,omitempty"` + MemAvg float64 `protobuf:"fixed64,4,opt,name=memAvg,proto3" json:"memAvg,omitempty"` + MemPeak float64 `protobuf:"fixed64,5,opt,name=memPeak,proto3" json:"memPeak,omitempty"` + // share of the snapshots the program was in the top lists + SeenPct float64 `protobuf:"fixed64,6,opt,name=seenPct,proto3" json:"seenPct,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ProcessUsage) Reset() { + *x = ProcessUsage{} + mi := &file_api_api_proto_msgTypes[23] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ProcessUsage) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ProcessUsage) ProtoMessage() {} + +func (x *ProcessUsage) ProtoReflect() protoreflect.Message { + mi := &file_api_api_proto_msgTypes[23] + 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 ProcessUsage.ProtoReflect.Descriptor instead. +func (*ProcessUsage) Descriptor() ([]byte, []int) { + return file_api_api_proto_rawDescGZIP(), []int{23} +} + +func (x *ProcessUsage) GetName() string { + if x != nil { + return x.Name + } + return "" +} + +func (x *ProcessUsage) GetCpuAvg() float64 { + if x != nil { + return x.CpuAvg + } + return 0 +} + +func (x *ProcessUsage) GetCpuPeak() float64 { + if x != nil { + return x.CpuPeak + } + return 0 +} + +func (x *ProcessUsage) GetMemAvg() float64 { + if x != nil { + return x.MemAvg + } + return 0 +} + +func (x *ProcessUsage) GetMemPeak() float64 { + if x != nil { + return x.MemPeak + } + return 0 +} + +func (x *ProcessUsage) GetSeenPct() float64 { + if x != nil { + return x.SeenPct + } + return 0 +} + +type ProcessUsageList struct { + state protoimpl.MessageState `protogen:"open.v1"` + Snapshots int32 `protobuf:"varint,1,opt,name=snapshots,proto3" json:"snapshots,omitempty"` + // time of the first snapshot in the range + FirstTime int64 `protobuf:"varint,2,opt,name=firstTime,proto3" json:"firstTime,omitempty"` + Processes []*ProcessUsage `protobuf:"bytes,3,rep,name=processes,proto3" json:"processes,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ProcessUsageList) Reset() { + *x = ProcessUsageList{} + mi := &file_api_api_proto_msgTypes[24] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ProcessUsageList) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ProcessUsageList) ProtoMessage() {} + +func (x *ProcessUsageList) ProtoReflect() protoreflect.Message { + mi := &file_api_api_proto_msgTypes[24] + 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 ProcessUsageList.ProtoReflect.Descriptor instead. +func (*ProcessUsageList) Descriptor() ([]byte, []int) { + return file_api_api_proto_rawDescGZIP(), []int{24} +} + +func (x *ProcessUsageList) GetSnapshots() int32 { + if x != nil { + return x.Snapshots + } + return 0 +} + +func (x *ProcessUsageList) GetFirstTime() int64 { + if x != nil { + return x.FirstTime + } + return 0 +} + +func (x *ProcessUsageList) GetProcesses() []*ProcessUsage { + if x != nil { + return x.Processes + } + return nil +} + +type SnapshotList struct { + state protoimpl.MessageState `protogen:"open.v1"` + Hosts []*HostSnapshot `protobuf:"bytes,1,rep,name=hosts,proto3" json:"hosts,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *SnapshotList) Reset() { + *x = SnapshotList{} + mi := &file_api_api_proto_msgTypes[25] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *SnapshotList) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*SnapshotList) ProtoMessage() {} + +func (x *SnapshotList) ProtoReflect() protoreflect.Message { + mi := &file_api_api_proto_msgTypes[25] + 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 SnapshotList.ProtoReflect.Descriptor instead. +func (*SnapshotList) Descriptor() ([]byte, []int) { + return file_api_api_proto_rawDescGZIP(), []int{25} +} + +func (x *SnapshotList) GetHosts() []*HostSnapshot { + if x != nil { + return x.Hosts + } + return nil +} + var File_api_api_proto protoreflect.FileDescriptor const file_api_api_proto_rawDesc = "" + @@ -1303,7 +1694,7 @@ const file_api_api_proto_rawDesc = "" + "\btimezone\x18\x03 \x01(\tR\btimezone\"H\n" + "\x0eEnrollResponse\x12\x1a\n" + "\bhostName\x18\x01 \x01(\tR\bhostName\x12\x1a\n" + - "\bagentKey\x18\x02 \x01(\tR\bagentKey\"\xa9\x03\n" + + "\bagentKey\x18\x02 \x01(\tR\bagentKey\"\xe3\x03\n" + "\vHostSummary\x12\x12\n" + "\x04name\x18\x01 \x01(\tR\x04name\x12\x0e\n" + "\x02up\x18\x02 \x01(\bR\x02up\x12\x1a\n" + @@ -1324,7 +1715,9 @@ const file_api_api_proto_rawDesc = "" + "\rworstSeverity\x18\x0e \x01(\x05R\rworstSeverity\x12\x1e\n" + "\n" + "containers\x18\x0f \x01(\x05R\n" + - "containers\"6\n" + + "containers\x12'\n" + + "\fdiskFullDays\x18\x10 \x01(\x01H\x00R\fdiskFullDays\x88\x01\x01B\x0f\n" + + "\r_diskFullDays\"6\n" + "\fFleetSummary\x12&\n" + "\x05hosts\x18\x01 \x03(\v2\x10.api.HostSummaryR\x05hosts\"!\n" + "\vHostRequest\x12\x12\n" + @@ -1383,7 +1776,36 @@ const file_api_api_proto_rawDesc = "" + " \x01(\x03R\n" + "resolvedAt\"5\n" + "\tAlertList\x12(\n" + - "\x06alerts\x18\x01 \x03(\v2\x10.api.AlertRecordR\x06alerts2\xd6\x04\n" + + "\x06alerts\x18\x01 \x03(\v2\x10.api.AlertRecordR\x06alerts\"\xca\x01\n" + + "\fDiskForecast\x12\x16\n" + + "\x06device\x18\x01 \x01(\tR\x06device\x12\x14\n" + + "\x05mount\x18\x02 \x01(\tR\x05mount\x12\x18\n" + + "\ausedPct\x18\x03 \x01(\x01R\ausedPct\x12\x1c\n" + + "\tpctPerDay\x18\x04 \x01(\x01R\tpctPerDay\x12 \n" + + "\vbytesPerDay\x18\x05 \x01(\x01R\vbytesPerDay\x12#\n" + + "\n" + + "daysToFull\x18\x06 \x01(\x01H\x00R\n" + + "daysToFull\x88\x01\x01B\r\n" + + "\v_daysToFull\";\n" + + "\x10DiskForecastList\x12'\n" + + "\x05disks\x18\x01 \x03(\v2\x11.api.DiskForecastR\x05disks\"M\n" + + "\x13ProcessUsageRequest\x12\x12\n" + + "\x04host\x18\x01 \x01(\tR\x04host\x12\x12\n" + + "\x04from\x18\x02 \x01(\x03R\x04from\x12\x0e\n" + + "\x02to\x18\x03 \x01(\x03R\x02to\"\xa0\x01\n" + + "\fProcessUsage\x12\x12\n" + + "\x04name\x18\x01 \x01(\tR\x04name\x12\x16\n" + + "\x06cpuAvg\x18\x02 \x01(\x01R\x06cpuAvg\x12\x18\n" + + "\acpuPeak\x18\x03 \x01(\x01R\acpuPeak\x12\x16\n" + + "\x06memAvg\x18\x04 \x01(\x01R\x06memAvg\x12\x18\n" + + "\amemPeak\x18\x05 \x01(\x01R\amemPeak\x12\x18\n" + + "\aseenPct\x18\x06 \x01(\x01R\aseenPct\"\x7f\n" + + "\x10ProcessUsageList\x12\x1c\n" + + "\tsnapshots\x18\x01 \x01(\x05R\tsnapshots\x12\x1c\n" + + "\tfirstTime\x18\x02 \x01(\x03R\tfirstTime\x12/\n" + + "\tprocesses\x18\x03 \x03(\v2\x11.api.ProcessUsageR\tprocesses\"7\n" + + "\fSnapshotList\x12'\n" + + "\x05hosts\x18\x01 \x03(\v2\x11.api.HostSnapshotR\x05hosts2\x82\x06\n" + "\x12MonitorDataService\x123\n" + "\x06Enroll\x12\x12.api.EnrollRequest\x1a\x13.api.EnrollResponse\"\x00\x12-\n" + "\n" + @@ -1396,7 +1818,10 @@ const file_api_api_proto_rawDesc = "" + "\vQuerySeries\x12\x12.api.SeriesRequest\x1a\x13.api.SeriesResponse\"\x00\x12<\n" + "\tProcesses\x12\x15.api.ProcessesRequest\x1a\x16.api.ProcessesResponse\"\x00\x126\n" + "\x11CustomMetricNames\x12\x10.api.HostRequest\x1a\r.api.NameList\"\x00\x12.\n" + - "\x06Alerts\x12\x12.api.AlertsRequest\x1a\x0e.api.AlertList\"\x00B)Z'github.com/dhamith93/SyMon/internal/apib\x06proto3" + "\x06Alerts\x12\x12.api.AlertsRequest\x1a\x0e.api.AlertList\"\x00\x12:\n" + + "\rDiskForecasts\x12\x10.api.HostRequest\x1a\x15.api.DiskForecastList\"\x00\x12A\n" + + "\fProcessUsage\x12\x18.api.ProcessUsageRequest\x1a\x15.api.ProcessUsageList\"\x00\x12+\n" + + "\tSnapshots\x12\t.api.Void\x1a\x11.api.SnapshotList\"\x00B)Z'github.com/dhamith93/SyMon/internal/apib\x06proto3" var ( file_api_api_proto_rawDescOnce sync.Once @@ -1410,61 +1835,76 @@ func file_api_api_proto_rawDescGZIP() []byte { return file_api_api_proto_rawDescData } -var file_api_api_proto_msgTypes = make([]protoimpl.MessageInfo, 20) +var file_api_api_proto_msgTypes = make([]protoimpl.MessageInfo, 26) var file_api_api_proto_goTypes = []any{ - (*Void)(nil), // 0: api.Void - (*Message)(nil), // 1: api.Message - (*ServerInfo)(nil), // 2: api.ServerInfo - (*MonitorData)(nil), // 3: api.MonitorData - (*EnrollRequest)(nil), // 4: api.EnrollRequest - (*EnrollResponse)(nil), // 5: api.EnrollResponse - (*HostSummary)(nil), // 6: api.HostSummary - (*FleetSummary)(nil), // 7: api.FleetSummary - (*HostRequest)(nil), // 8: api.HostRequest - (*HostSnapshot)(nil), // 9: api.HostSnapshot - (*SeriesRequest)(nil), // 10: api.SeriesRequest - (*Point)(nil), // 11: api.Point - (*Series)(nil), // 12: api.Series - (*SeriesResponse)(nil), // 13: api.SeriesResponse - (*ProcessesRequest)(nil), // 14: api.ProcessesRequest - (*ProcessesResponse)(nil), // 15: api.ProcessesResponse - (*NameList)(nil), // 16: api.NameList - (*AlertsRequest)(nil), // 17: api.AlertsRequest - (*AlertRecord)(nil), // 18: api.AlertRecord - (*AlertList)(nil), // 19: api.AlertList + (*Void)(nil), // 0: api.Void + (*Message)(nil), // 1: api.Message + (*ServerInfo)(nil), // 2: api.ServerInfo + (*MonitorData)(nil), // 3: api.MonitorData + (*EnrollRequest)(nil), // 4: api.EnrollRequest + (*EnrollResponse)(nil), // 5: api.EnrollResponse + (*HostSummary)(nil), // 6: api.HostSummary + (*FleetSummary)(nil), // 7: api.FleetSummary + (*HostRequest)(nil), // 8: api.HostRequest + (*HostSnapshot)(nil), // 9: api.HostSnapshot + (*SeriesRequest)(nil), // 10: api.SeriesRequest + (*Point)(nil), // 11: api.Point + (*Series)(nil), // 12: api.Series + (*SeriesResponse)(nil), // 13: api.SeriesResponse + (*ProcessesRequest)(nil), // 14: api.ProcessesRequest + (*ProcessesResponse)(nil), // 15: api.ProcessesResponse + (*NameList)(nil), // 16: api.NameList + (*AlertsRequest)(nil), // 17: api.AlertsRequest + (*AlertRecord)(nil), // 18: api.AlertRecord + (*AlertList)(nil), // 19: api.AlertList + (*DiskForecast)(nil), // 20: api.DiskForecast + (*DiskForecastList)(nil), // 21: api.DiskForecastList + (*ProcessUsageRequest)(nil), // 22: api.ProcessUsageRequest + (*ProcessUsage)(nil), // 23: api.ProcessUsage + (*ProcessUsageList)(nil), // 24: api.ProcessUsageList + (*SnapshotList)(nil), // 25: api.SnapshotList } var file_api_api_proto_depIdxs = []int32{ 6, // 0: api.FleetSummary.hosts:type_name -> api.HostSummary 11, // 1: api.Series.points:type_name -> api.Point 12, // 2: api.SeriesResponse.series:type_name -> api.Series 18, // 3: api.AlertList.alerts:type_name -> api.AlertRecord - 4, // 4: api.MonitorDataService.Enroll:input_type -> api.EnrollRequest - 2, // 5: api.MonitorDataService.HandlePing:input_type -> api.ServerInfo - 2, // 6: api.MonitorDataService.InitAgent:input_type -> api.ServerInfo - 3, // 7: api.MonitorDataService.HandleMonitorData:input_type -> api.MonitorData - 3, // 8: api.MonitorDataService.HandleCustomMonitorData:input_type -> api.MonitorData - 0, // 9: api.MonitorDataService.Fleet:input_type -> api.Void - 8, // 10: api.MonitorDataService.Snapshot:input_type -> api.HostRequest - 10, // 11: api.MonitorDataService.QuerySeries:input_type -> api.SeriesRequest - 14, // 12: api.MonitorDataService.Processes:input_type -> api.ProcessesRequest - 8, // 13: api.MonitorDataService.CustomMetricNames:input_type -> api.HostRequest - 17, // 14: api.MonitorDataService.Alerts:input_type -> api.AlertsRequest - 5, // 15: api.MonitorDataService.Enroll:output_type -> api.EnrollResponse - 1, // 16: api.MonitorDataService.HandlePing:output_type -> api.Message - 1, // 17: api.MonitorDataService.InitAgent:output_type -> api.Message - 1, // 18: api.MonitorDataService.HandleMonitorData:output_type -> api.Message - 1, // 19: api.MonitorDataService.HandleCustomMonitorData:output_type -> api.Message - 7, // 20: api.MonitorDataService.Fleet:output_type -> api.FleetSummary - 9, // 21: api.MonitorDataService.Snapshot:output_type -> api.HostSnapshot - 13, // 22: api.MonitorDataService.QuerySeries:output_type -> api.SeriesResponse - 15, // 23: api.MonitorDataService.Processes:output_type -> api.ProcessesResponse - 16, // 24: api.MonitorDataService.CustomMetricNames:output_type -> api.NameList - 19, // 25: api.MonitorDataService.Alerts:output_type -> api.AlertList - 15, // [15:26] is the sub-list for method output_type - 4, // [4:15] is the sub-list for method input_type - 4, // [4:4] is the sub-list for extension type_name - 4, // [4:4] is the sub-list for extension extendee - 0, // [0:4] is the sub-list for field type_name + 20, // 4: api.DiskForecastList.disks:type_name -> api.DiskForecast + 23, // 5: api.ProcessUsageList.processes:type_name -> api.ProcessUsage + 9, // 6: api.SnapshotList.hosts:type_name -> api.HostSnapshot + 4, // 7: api.MonitorDataService.Enroll:input_type -> api.EnrollRequest + 2, // 8: api.MonitorDataService.HandlePing:input_type -> api.ServerInfo + 2, // 9: api.MonitorDataService.InitAgent:input_type -> api.ServerInfo + 3, // 10: api.MonitorDataService.HandleMonitorData:input_type -> api.MonitorData + 3, // 11: api.MonitorDataService.HandleCustomMonitorData:input_type -> api.MonitorData + 0, // 12: api.MonitorDataService.Fleet:input_type -> api.Void + 8, // 13: api.MonitorDataService.Snapshot:input_type -> api.HostRequest + 10, // 14: api.MonitorDataService.QuerySeries:input_type -> api.SeriesRequest + 14, // 15: api.MonitorDataService.Processes:input_type -> api.ProcessesRequest + 8, // 16: api.MonitorDataService.CustomMetricNames:input_type -> api.HostRequest + 17, // 17: api.MonitorDataService.Alerts:input_type -> api.AlertsRequest + 8, // 18: api.MonitorDataService.DiskForecasts:input_type -> api.HostRequest + 22, // 19: api.MonitorDataService.ProcessUsage:input_type -> api.ProcessUsageRequest + 0, // 20: api.MonitorDataService.Snapshots:input_type -> api.Void + 5, // 21: api.MonitorDataService.Enroll:output_type -> api.EnrollResponse + 1, // 22: api.MonitorDataService.HandlePing:output_type -> api.Message + 1, // 23: api.MonitorDataService.InitAgent:output_type -> api.Message + 1, // 24: api.MonitorDataService.HandleMonitorData:output_type -> api.Message + 1, // 25: api.MonitorDataService.HandleCustomMonitorData:output_type -> api.Message + 7, // 26: api.MonitorDataService.Fleet:output_type -> api.FleetSummary + 9, // 27: api.MonitorDataService.Snapshot:output_type -> api.HostSnapshot + 13, // 28: api.MonitorDataService.QuerySeries:output_type -> api.SeriesResponse + 15, // 29: api.MonitorDataService.Processes:output_type -> api.ProcessesResponse + 16, // 30: api.MonitorDataService.CustomMetricNames:output_type -> api.NameList + 19, // 31: api.MonitorDataService.Alerts:output_type -> api.AlertList + 21, // 32: api.MonitorDataService.DiskForecasts:output_type -> api.DiskForecastList + 24, // 33: api.MonitorDataService.ProcessUsage:output_type -> api.ProcessUsageList + 25, // 34: api.MonitorDataService.Snapshots:output_type -> api.SnapshotList + 21, // [21:35] is the sub-list for method output_type + 7, // [7:21] is the sub-list for method input_type + 7, // [7:7] is the sub-list for extension type_name + 7, // [7:7] is the sub-list for extension extendee + 0, // [0:7] is the sub-list for field type_name } func init() { file_api_api_proto_init() } @@ -1472,13 +1912,15 @@ func file_api_api_proto_init() { if File_api_api_proto != nil { return } + file_api_api_proto_msgTypes[6].OneofWrappers = []any{} + file_api_api_proto_msgTypes[20].OneofWrappers = []any{} type x struct{} out := protoimpl.TypeBuilder{ File: protoimpl.DescBuilder{ GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_api_api_proto_rawDesc), len(file_api_api_proto_rawDesc)), NumEnums: 0, - NumMessages: 20, + NumMessages: 26, NumExtensions: 0, NumServices: 1, }, diff --git a/internal/api/api.proto b/internal/api/api.proto index b90af8a..c8f283d 100644 --- a/internal/api/api.proto +++ b/internal/api/api.proto @@ -52,6 +52,8 @@ message HostSummary { int32 worstSeverity = 14; // running containers at the latest snapshot int32 containers = 15; + // days until the first disk fills up, unset when none is filling up + optional double diskFullDays = 16; } message FleetSummary { @@ -144,6 +146,50 @@ message AlertList { repeated AlertRecord alerts = 1; } +// A disk's growth over the last week and when it fills up at that rate +message DiskForecast { + string device = 1; + string mount = 2; + double usedPct = 3; + double pctPerDay = 4; + double bytesPerDay = 5; + // unset when the disk is not filling up + optional double daysToFull = 6; +} + +message DiskForecastList { + repeated DiskForecast disks = 1; +} + +message ProcessUsageRequest { + string host = 1; + int64 from = 2; + int64 to = 3; +} + +// One program's share of the range. Processes with the same name are +// added up per snapshot. +message ProcessUsage { + string name = 1; + double cpuAvg = 2; + double cpuPeak = 3; + double memAvg = 4; + double memPeak = 5; + // share of the snapshots the program was in the top lists + double seenPct = 6; +} + +message ProcessUsageList { + int32 snapshots = 1; + // time of the first snapshot in the range + int64 firstTime = 2; + repeated ProcessUsage processes = 3; +} + +message SnapshotList { + repeated HostSnapshot hosts = 1; +} + service MonitorDataService { // agent rpc Enroll(EnrollRequest) returns (EnrollResponse) {} @@ -159,4 +205,8 @@ service MonitorDataService { rpc Processes(ProcessesRequest) returns (ProcessesResponse) {} rpc CustomMetricNames(HostRequest) returns (NameList) {} rpc Alerts(AlertsRequest) returns (AlertList) {} + rpc DiskForecasts(HostRequest) returns (DiskForecastList) {} + rpc ProcessUsage(ProcessUsageRequest) returns (ProcessUsageList) {} + // every host's latest snapshot, for the metrics endpoint + rpc Snapshots(Void) returns (SnapshotList) {} } diff --git a/internal/api/api_grpc.pb.go b/internal/api/api_grpc.pb.go index 5df5e5f..80b5f5f 100644 --- a/internal/api/api_grpc.pb.go +++ b/internal/api/api_grpc.pb.go @@ -30,6 +30,9 @@ const ( MonitorDataService_Processes_FullMethodName = "/api.MonitorDataService/Processes" MonitorDataService_CustomMetricNames_FullMethodName = "/api.MonitorDataService/CustomMetricNames" MonitorDataService_Alerts_FullMethodName = "/api.MonitorDataService/Alerts" + MonitorDataService_DiskForecasts_FullMethodName = "/api.MonitorDataService/DiskForecasts" + MonitorDataService_ProcessUsage_FullMethodName = "/api.MonitorDataService/ProcessUsage" + MonitorDataService_Snapshots_FullMethodName = "/api.MonitorDataService/Snapshots" ) // MonitorDataServiceClient is the client API for MonitorDataService service. @@ -49,6 +52,10 @@ type MonitorDataServiceClient interface { Processes(ctx context.Context, in *ProcessesRequest, opts ...grpc.CallOption) (*ProcessesResponse, error) CustomMetricNames(ctx context.Context, in *HostRequest, opts ...grpc.CallOption) (*NameList, error) Alerts(ctx context.Context, in *AlertsRequest, opts ...grpc.CallOption) (*AlertList, error) + DiskForecasts(ctx context.Context, in *HostRequest, opts ...grpc.CallOption) (*DiskForecastList, error) + ProcessUsage(ctx context.Context, in *ProcessUsageRequest, opts ...grpc.CallOption) (*ProcessUsageList, error) + // every host's latest snapshot, for the metrics endpoint + Snapshots(ctx context.Context, in *Void, opts ...grpc.CallOption) (*SnapshotList, error) } type monitorDataServiceClient struct { @@ -169,6 +176,36 @@ func (c *monitorDataServiceClient) Alerts(ctx context.Context, in *AlertsRequest return out, nil } +func (c *monitorDataServiceClient) DiskForecasts(ctx context.Context, in *HostRequest, opts ...grpc.CallOption) (*DiskForecastList, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(DiskForecastList) + err := c.cc.Invoke(ctx, MonitorDataService_DiskForecasts_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *monitorDataServiceClient) ProcessUsage(ctx context.Context, in *ProcessUsageRequest, opts ...grpc.CallOption) (*ProcessUsageList, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ProcessUsageList) + err := c.cc.Invoke(ctx, MonitorDataService_ProcessUsage_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *monitorDataServiceClient) Snapshots(ctx context.Context, in *Void, opts ...grpc.CallOption) (*SnapshotList, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(SnapshotList) + err := c.cc.Invoke(ctx, MonitorDataService_Snapshots_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + // MonitorDataServiceServer is the server API for MonitorDataService service. // All implementations must embed UnimplementedMonitorDataServiceServer // for forward compatibility. @@ -186,6 +223,10 @@ type MonitorDataServiceServer interface { Processes(context.Context, *ProcessesRequest) (*ProcessesResponse, error) CustomMetricNames(context.Context, *HostRequest) (*NameList, error) Alerts(context.Context, *AlertsRequest) (*AlertList, error) + DiskForecasts(context.Context, *HostRequest) (*DiskForecastList, error) + ProcessUsage(context.Context, *ProcessUsageRequest) (*ProcessUsageList, error) + // every host's latest snapshot, for the metrics endpoint + Snapshots(context.Context, *Void) (*SnapshotList, error) mustEmbedUnimplementedMonitorDataServiceServer() } @@ -229,6 +270,15 @@ func (UnimplementedMonitorDataServiceServer) CustomMetricNames(context.Context, func (UnimplementedMonitorDataServiceServer) Alerts(context.Context, *AlertsRequest) (*AlertList, error) { return nil, status.Errorf(codes.Unimplemented, "method Alerts not implemented") } +func (UnimplementedMonitorDataServiceServer) DiskForecasts(context.Context, *HostRequest) (*DiskForecastList, error) { + return nil, status.Errorf(codes.Unimplemented, "method DiskForecasts not implemented") +} +func (UnimplementedMonitorDataServiceServer) ProcessUsage(context.Context, *ProcessUsageRequest) (*ProcessUsageList, error) { + return nil, status.Errorf(codes.Unimplemented, "method ProcessUsage not implemented") +} +func (UnimplementedMonitorDataServiceServer) Snapshots(context.Context, *Void) (*SnapshotList, error) { + return nil, status.Errorf(codes.Unimplemented, "method Snapshots not implemented") +} func (UnimplementedMonitorDataServiceServer) mustEmbedUnimplementedMonitorDataServiceServer() {} func (UnimplementedMonitorDataServiceServer) testEmbeddedByValue() {} @@ -448,6 +498,60 @@ func _MonitorDataService_Alerts_Handler(srv interface{}, ctx context.Context, de return interceptor(ctx, in, info, handler) } +func _MonitorDataService_DiskForecasts_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(HostRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(MonitorDataServiceServer).DiskForecasts(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: MonitorDataService_DiskForecasts_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(MonitorDataServiceServer).DiskForecasts(ctx, req.(*HostRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _MonitorDataService_ProcessUsage_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ProcessUsageRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(MonitorDataServiceServer).ProcessUsage(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: MonitorDataService_ProcessUsage_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(MonitorDataServiceServer).ProcessUsage(ctx, req.(*ProcessUsageRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _MonitorDataService_Snapshots_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(Void) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(MonitorDataServiceServer).Snapshots(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: MonitorDataService_Snapshots_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(MonitorDataServiceServer).Snapshots(ctx, req.(*Void)) + } + return interceptor(ctx, in, info, handler) +} + // MonitorDataService_ServiceDesc is the grpc.ServiceDesc for MonitorDataService service. // It's only intended for direct use with grpc.RegisterService, // and not to be introspected or modified (even as a copy) @@ -499,6 +603,18 @@ var MonitorDataService_ServiceDesc = grpc.ServiceDesc{ MethodName: "Alerts", Handler: _MonitorDataService_Alerts_Handler, }, + { + MethodName: "DiskForecasts", + Handler: _MonitorDataService_DiskForecasts_Handler, + }, + { + MethodName: "ProcessUsage", + Handler: _MonitorDataService_ProcessUsage_Handler, + }, + { + MethodName: "Snapshots", + Handler: _MonitorDataService_Snapshots_Handler, + }, }, Streams: []grpc.StreamDesc{}, Metadata: "api/api.proto", From c3e1098a5439e30f8431a63ca584d76910005005 Mon Sep 17 00:00:00 2001 From: Dhamith Hewamullage Date: Tue, 29 Sep 2026 19:21:36 +0530 Subject: [PATCH 02/10] forecast when disks fill up --- client/internal/server/server.go | 35 ++++++ client/internal/server/server_test.go | 32 +++++- client/web/src/lib/api.ts | 17 +++ client/web/src/lib/format.test.ts | 19 +++- client/web/src/lib/format.ts | 11 ++ client/web/src/pages/Fleet.svelte | 29 +++-- client/web/src/pages/Host.svelte | 31 +++++- internal/api/api.go | 57 +++++++++- internal/store/forecast.go | 154 ++++++++++++++++++++++++++ internal/store/forecast_test.go | 123 ++++++++++++++++++++ 10 files changed, 493 insertions(+), 15 deletions(-) create mode 100644 internal/store/forecast.go create mode 100644 internal/store/forecast_test.go diff --git a/client/internal/server/server.go b/client/internal/server/server.go index ca0a9c9..8ea9e17 100644 --- a/client/internal/server/server.go +++ b/client/internal/server/server.go @@ -68,6 +68,7 @@ func (s *server) routes() http.Handler { mux.HandleFunc("GET /api/v1/hosts/{host}/series", s.getSeries) mux.HandleFunc("GET /api/v1/hosts/{host}/processes", s.getProcesses) mux.HandleFunc("GET /api/v1/hosts/{host}/custom-metrics", s.getCustomMetrics) + mux.HandleFunc("GET /api/v1/hosts/{host}/disk-forecasts", s.getDiskForecasts) mux.HandleFunc("GET /api/v1/alerts", s.getAlerts) mux.HandleFunc("/api/", func(w http.ResponseWriter, r *http.Request) { writeError(w, http.StatusNotFound, "no such endpoint") @@ -98,6 +99,8 @@ type hostSummary struct { ActiveAlerts int32 `json:"activeAlerts"` WorstSeverity int32 `json:"worstSeverity"` Containers int32 `json:"containers"` + // null when no disk is filling up + DiskFullDays *float64 `json:"diskFullDays"` } func (s *server) getFleet(w http.ResponseWriter, r *http.Request) { @@ -124,6 +127,7 @@ func (s *server) getFleet(w http.ResponseWriter, r *http.Request) { ActiveAlerts: h.ActiveAlerts, WorstSeverity: h.WorstSeverity, Containers: h.Containers, + DiskFullDays: h.DiskFullDays, }) } writeJSON(w, map[string]any{"hosts": hosts}) @@ -216,6 +220,37 @@ func (s *server) getCustomMetrics(w http.ResponseWriter, r *http.Request) { writeJSON(w, map[string]any{"names": nonNil(names.Names)}) } +type diskForecast struct { + Device string `json:"device"` + Mount string `json:"mount"` + UsedPct float64 `json:"usedPct"` + PctPerDay float64 `json:"pctPerDay"` + BytesPerDay float64 `json:"bytesPerDay"` + // null when the disk is not filling up + DaysToFull *float64 `json:"daysToFull"` +} + +func (s *server) getDiskForecasts(w http.ResponseWriter, r *http.Request) { + host := r.PathValue("host") + response, err := s.collector.DiskForecasts(r.Context(), &api.HostRequest{Host: host}) + if err != nil { + writeGRPCError(w, "disk forecasts of "+host, err) + return + } + disks := make([]diskForecast, 0, len(response.Disks)) + for _, d := range response.Disks { + disks = append(disks, diskForecast{ + Device: d.Device, + Mount: d.Mount, + UsedPct: d.UsedPct, + PctPerDay: d.PctPerDay, + BytesPerDay: d.BytesPerDay, + DaysToFull: d.DaysToFull, + }) + } + writeJSON(w, map[string]any{"disks": disks}) +} + type alert struct { ID int64 `json:"id"` Host string `json:"host"` diff --git a/client/internal/server/server_test.go b/client/internal/server/server_test.go index 7b2389d..7b31f42 100644 --- a/client/internal/server/server_test.go +++ b/client/internal/server/server_test.go @@ -26,7 +26,10 @@ type fakeCollector struct { } func (f *fakeCollector) Fleet(ctx context.Context, in *api.Void) (*api.FleetSummary, error) { - return &api.FleetSummary{Hosts: []*api.HostSummary{{Name: "web1", Up: true, CpuPct: 37, ActiveAlerts: 2}}}, nil + return &api.FleetSummary{Hosts: []*api.HostSummary{ + {Name: "web1", Up: true, CpuPct: 37, ActiveAlerts: 2, DiskFullDays: floatPtr(12.5)}, + {Name: "db1", Up: true}, + }}, nil } func (f *fakeCollector) Snapshot(ctx context.Context, in *api.HostRequest) (*api.HostSnapshot, error) { @@ -55,6 +58,17 @@ func (f *fakeCollector) CustomMetricNames(ctx context.Context, in *api.HostReque return &api.NameList{}, nil } +func (f *fakeCollector) DiskForecasts(ctx context.Context, in *api.HostRequest) (*api.DiskForecastList, error) { + return &api.DiskForecastList{Disks: []*api.DiskForecast{ + {Device: "/dev/sda1", Mount: "/", UsedPct: 40}, + {Device: "/dev/sdb1", Mount: "/data", UsedPct: 60, PctPerDay: 2, BytesPerDay: 2e7, DaysToFull: floatPtr(20)}, + }}, nil +} + +func floatPtr(v float64) *float64 { + return &v +} + func (f *fakeCollector) Alerts(ctx context.Context, in *api.AlertsRequest) (*api.AlertList, error) { return &api.AlertList{Alerts: []*api.AlertRecord{{Id: 7, Host: "web1", Rule: "CPU", Severity: 2, StartedAt: 1700000000}}}, nil } @@ -106,7 +120,21 @@ func TestFleet(t *testing.T) { if err := json.Unmarshal([]byte(body), &out); err != nil { t.Fatal(err) } - if code != 200 || len(out.Hosts) != 1 || out.Hosts[0].Name != "web1" || out.Hosts[0].CPUPct != 37 || out.Hosts[0].ActiveAlerts != 2 { + if code != 200 || len(out.Hosts) != 2 || out.Hosts[0].Name != "web1" || out.Hosts[0].CPUPct != 37 || out.Hosts[0].ActiveAlerts != 2 { + t.Errorf("unexpected response %d: %s", code, body) + } + if !strings.Contains(body, `"diskFullDays":12.5`) || !strings.Contains(body, `"diskFullDays":null`) { + t.Errorf("expected a forecast for web1 and null for db1: %s", body) + } +} + +func TestDiskForecasts(t *testing.T) { + s, _ := newTestServer(t, nil) + code, body, _ := get(t, s, "/api/v1/hosts/web1/disk-forecasts") + want := `{"disks":[` + + `{"device":"/dev/sda1","mount":"/","usedPct":40,"pctPerDay":0,"bytesPerDay":0,"daysToFull":null},` + + `{"device":"/dev/sdb1","mount":"/data","usedPct":60,"pctPerDay":2,"bytesPerDay":20000000,"daysToFull":20}]}` + if code != 200 || strings.TrimSpace(body) != want { t.Errorf("unexpected response %d: %s", code, body) } } diff --git a/client/web/src/lib/api.ts b/client/web/src/lib/api.ts index 79dd1a9..9fb7f0e 100644 --- a/client/web/src/lib/api.ts +++ b/client/web/src/lib/api.ts @@ -18,8 +18,24 @@ export interface HostSummary { worstSeverity: number; // running containers at the latest snapshot containers: number; + // days until the first disk is full, null when none is filling up + diskFullDays: number | null; } +// A disk's growth over the last week and when it fills up at that rate +export interface DiskForecast { + device: string; + mount: string; + usedPct: number; + pctPerDay: number; + bytesPerDay: number; + // null when the disk is not filling up + daysToFull: number | null; +} + +// disks that fill up sooner than this many days are shown as warnings +export const diskFullSoonDays = 30; + export interface SeriesData { label: string; // [unix seconds, value] @@ -175,6 +191,7 @@ export const api = { get(`${host(name)}/series`, { metric, from, to, ...options }), processes: (name: string, at?: number) => get<{ time: number; processes: Processes }>(`${host(name)}/processes`, { at }), customMetrics: (name: string) => get<{ names: string[] }>(`${host(name)}/custom-metrics`), + diskForecasts: (name: string) => get<{ disks: DiskForecast[] }>(`${host(name)}/disk-forecasts`), alerts: (filter: { host?: string; open?: boolean; from?: number; to?: number } = {}) => get<{ alerts: AlertRecord[] }>('/api/v1/alerts', filter), }; diff --git a/client/web/src/lib/format.test.ts b/client/web/src/lib/format.test.ts index a63b02c..94debba 100644 --- a/client/web/src/lib/format.test.ts +++ b/client/web/src/lib/format.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it } from 'vitest'; -import { formatAgo, formatBytes, formatDuration, formatNumber, formatPercent, formatRate, formatValue } from './format'; +import { formatAgo, formatBytes, formatDays, formatDuration, formatNumber, formatPercent, formatRate, formatValue } from './format'; describe('formatBytes', () => { it('uses binary units', () => { @@ -45,3 +45,20 @@ describe('formatDuration', () => { expect(formatAgo(1000, 1000 + 180)).toBe('3m ago'); }); }); + +describe('formatDays', () => { + it('rounds to a unit that fits', () => { + expect(formatDays(0)).toBe('under a day'); + expect(formatDays(0.4)).toBe('under a day'); + expect(formatDays(1.2)).toBe('about 1 day'); + expect(formatDays(9.4)).toBe('about 9 days'); + expect(formatDays(20.4)).toBe('about 20 days'); + expect(formatDays(35)).toBe('about 5 weeks'); + expect(formatDays(120)).toBe('about 4 months'); + }); + + it('handles bad input', () => { + expect(formatDays(-1)).toBe('–'); + expect(formatDays(NaN)).toBe('–'); + }); +}); diff --git a/client/web/src/lib/format.ts b/client/web/src/lib/format.ts index fae3d28..27b27ee 100644 --- a/client/web/src/lib/format.ts +++ b/client/web/src/lib/format.ts @@ -64,6 +64,17 @@ export function formatAgo(unixSeconds: number, now = Date.now() / 1000): string return `${formatDuration(diff)} ago`; } +// how long until something happens: "under a day", "about 9 days", +// "about 5 weeks", "about 4 months" +export function formatDays(days: number): string { + if (!Number.isFinite(days) || days < 0) return '–'; + if (days < 1) return 'under a day'; + if (days < 1.5) return 'about 1 day'; + if (days < 21) return `about ${Math.round(days)} days`; + if (days < 60) return `about ${Math.round(days / 7)} weeks`; + return `about ${Math.round(days / 30)} months`; +} + export function formatDateTime(unixSeconds: number): string { if (!unixSeconds) return '–'; return new Date(unixSeconds * 1000).toLocaleString(undefined, { diff --git a/client/web/src/pages/Fleet.svelte b/client/web/src/pages/Fleet.svelte index 44ede56..bdfb9f3 100644 --- a/client/web/src/pages/Fleet.svelte +++ b/client/web/src/pages/Fleet.svelte @@ -1,8 +1,8 @@ + +
+
+

Busiest over this range

+
+ + +
+
+

+ {#if snapshots > 0} + From {snapshots.toLocaleString()} snapshots since {formatDateTime(firstTime)}. Each keeps only the top 10 processes, so + averages are a lower bound. + {:else if !error && !loading} + No process lists recorded in this range. + {/if} +

+ + {#if error} +

{error}

+ {:else if rows.length > 0} +
+ + + + + + + + + + + {#each rows as process (process.name)} + + + + + + + {/each} + +
ProgramAveragePeakSeen
{process.name}{formatPercent(view === 'CPU' ? process.cpuAvg : process.memAvg, 1)}{formatPercent(view === 'CPU' ? process.cpuPeak : process.memPeak, 1)}{formatPercent(process.seenPct)}
+
+ {/if} +
+ + diff --git a/client/web/src/lib/api.ts b/client/web/src/lib/api.ts index 9fb7f0e..912e4d0 100644 --- a/client/web/src/lib/api.ts +++ b/client/web/src/lib/api.ts @@ -76,6 +76,17 @@ export interface Process { Threads: number; } +// One program's share of a time range, with its processes added up +export interface ProcessUsage { + name: string; + cpuAvg: number; + cpuPeak: number; + memAvg: number; + memPeak: number; + // share of snapshots the program was in the top lists + seenPct: number; +} + export interface Processes { CPU: Process[] | null; Memory: Process[] | null; @@ -190,6 +201,8 @@ export const api = { series: (name: string, metric: string, from: number, to: number, options: { label?: string; maxPoints?: number; max?: boolean } = {}) => get(`${host(name)}/series`, { metric, from, to, ...options }), processes: (name: string, at?: number) => get<{ time: number; processes: Processes }>(`${host(name)}/processes`, { at }), + processUsage: (name: string, from: number, to: number) => + get<{ snapshots: number; firstTime: number; processes: ProcessUsage[] }>(`${host(name)}/process-usage`, { from, to }), customMetrics: (name: string) => get<{ names: string[] }>(`${host(name)}/custom-metrics`), diskForecasts: (name: string) => get<{ disks: DiskForecast[] }>(`${host(name)}/disk-forecasts`), alerts: (filter: { host?: string; open?: boolean; from?: number; to?: number } = {}) => diff --git a/client/web/src/pages/Host.svelte b/client/web/src/pages/Host.svelte index a7d1544..029519a 100644 --- a/client/web/src/pages/Host.svelte +++ b/client/web/src/pages/Host.svelte @@ -8,6 +8,7 @@ import { poll } from '../lib/poll'; import { hostPath, location, navigate } from '../lib/router.svelte'; import { rangeQuery, resolveRange } from '../lib/timerange'; + import BusiestProcesses from '../components/BusiestProcesses.svelte'; import ChartCard from '../components/ChartCard.svelte'; import ContainerTable from '../components/ContainerTable.svelte'; import Heatmap from '../components/Heatmap.svelte'; @@ -216,7 +217,10 @@

Processes

- (processesAt = 0)} /> +
+ (processesAt = 0)} /> + +
{#if snapshot} @@ -376,6 +380,13 @@ gap: 12px; } + .processes { + display: grid; + grid-template-columns: repeat(auto-fit, minmax(min(100%, 560px), 1fr)); + gap: 12px; + align-items: start; + } + .box { padding: 14px; min-width: 0; diff --git a/internal/api/api.go b/internal/api/api.go index 60932eb..3430fa2 100644 --- a/internal/api/api.go +++ b/internal/api/api.go @@ -258,6 +258,25 @@ func (s *Server) Processes(ctx context.Context, in *ProcessesRequest) (*Processe return &ProcessesResponse{Time: snapshotTime.Unix(), ProcessesJson: string(processes)}, nil } +func (s *Server) ProcessUsage(ctx context.Context, in *ProcessUsageRequest) (*ProcessUsageList, error) { + result, err := s.Store.ProcessUsage(ctx, in.Host, time.Unix(in.From, 0), time.Unix(in.To, 0)) + if err != nil { + return nil, toStatus(err) + } + list := &ProcessUsageList{Snapshots: int32(result.Snapshots), FirstTime: unix(result.FirstTime)} + for _, usage := range result.Processes { + list.Processes = append(list.Processes, &ProcessUsage{ + Name: usage.Name, + CpuAvg: usage.CPUAvg, + CpuPeak: usage.CPUPeak, + MemAvg: usage.MemAvg, + MemPeak: usage.MemPeak, + SeenPct: usage.SeenPct, + }) + } + return list, nil +} + func (s *Server) CustomMetricNames(ctx context.Context, in *HostRequest) (*NameList, error) { names, err := s.Store.CustomMetricNames(ctx, in.Host) if err != nil { diff --git a/internal/store/query.go b/internal/store/query.go index 29cdbe3..7a5bc80 100644 --- a/internal/store/query.go +++ b/internal/store/query.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "fmt" "time" "github.com/dhamith93/SyMon/internal/monitor" @@ -133,6 +134,100 @@ func (s *Store) Processes(ctx context.Context, host string, at time.Time) (time. return snapshotTime, processes, err } +// ProcessUsage is one program's share of a time range. Processes with the +// same name, like the workers of a web server, are added up per snapshot. +type ProcessUsage struct { + Name string + CPUAvg float64 + CPUPeak float64 + MemAvg float64 + MemPeak float64 + // SeenPct is the share of snapshots the program was in the top lists + SeenPct float64 +} + +type ProcessUsageResult struct { + Snapshots int + // FirstTime is the first snapshot in the range, zero when there is none + FirstTime time.Time + // Processes has the busiest programs by CPU and by memory, busiest CPU first + Processes []ProcessUsage +} + +// processUsageLimit is how many programs ProcessUsage returns for each of +// CPU and memory +const processUsageLimit = 15 + +// ProcessUsage adds up each program's CPU and memory over a range. Snapshots +// only keep the top processes, so a program counts as 0 where it was not +// among them, and the averages are a lower bound. +func (s *Store) ProcessUsage(ctx context.Context, host string, from time.Time, to time.Time) (ProcessUsageResult, error) { + if !to.After(from) { + return ProcessUsageResult{}, fmt.Errorf("%w: from must be before to", ErrInvalid) + } + hostID, err := s.hostID(ctx, host) + if err != nil { + return ProcessUsageResult{}, err + } + + result := ProcessUsageResult{Processes: []ProcessUsage{}} + var first *time.Time + err = s.pool.QueryRow(ctx, ` + SELECT count(*), min(time) FROM process_snapshots + WHERE host_id = $1 AND time >= $2 AND time < $3`, hostID, from, to).Scan(&result.Snapshots, &first) + if err != nil || result.Snapshots == 0 { + return result, err + } + result.FirstTime = *first + + // the CPU and Memory lists overlap, so each pid counts once per snapshot + rows, err := s.pool.Query(ctx, ` + WITH processes AS ( + SELECT DISTINCT ON (s.time, p->>'Pid') + s.time, + coalesce(nullif(p->>'Name', ''), nullif(p->>'ExecPath', ''), 'unknown') AS name, + coalesce((p->>'CPUUsage')::float8, 0) AS cpu, + coalesce((p->>'MemUsage')::float8, 0) AS mem + FROM process_snapshots s, + jsonb_path_query(s.processes, '$.*[*] ? (@.type() == "object")') AS p + WHERE s.host_id = $1 AND s.time >= $2 AND s.time < $3 + ), + per_snapshot AS ( + SELECT time, name, sum(cpu) AS cpu, sum(mem) AS mem + FROM processes + GROUP BY time, name + ), + ranked AS ( + SELECT name, sum(cpu) AS cpu_sum, max(cpu) AS cpu_peak, sum(mem) AS mem_sum, max(mem) AS mem_peak, count(*) AS seen, + row_number() OVER (ORDER BY sum(cpu) DESC, name) AS cpu_rank, + row_number() OVER (ORDER BY sum(mem) DESC, name) AS mem_rank + FROM per_snapshot + GROUP BY name + ) + SELECT name, cpu_sum, cpu_peak, mem_sum, mem_peak, seen FROM ranked + WHERE cpu_rank <= $4 OR mem_rank <= $4 + ORDER BY cpu_sum DESC, name`, hostID, from, to, processUsageLimit) + if err != nil { + return ProcessUsageResult{}, err + } + defer rows.Close() + + snapshots := float64(result.Snapshots) + for rows.Next() { + var usage ProcessUsage + var cpuSum, memSum float64 + var seen int + if err := rows.Scan(&usage.Name, &cpuSum, &usage.CPUPeak, &memSum, &usage.MemPeak, &seen); err != nil { + return ProcessUsageResult{}, err + } + usage.CPUAvg = cpuSum / snapshots + usage.MemAvg = memSum / snapshots + usage.SeenPct = 100 * float64(seen) / snapshots + result.Processes = append(result.Processes, usage) + } + return result, rows.Err() +} + // CustomMetricNames lists the custom metrics a host has sent within the // hourly retention func (s *Store) CustomMetricNames(ctx context.Context, host string) ([]string, error) { diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 1f7ae3a..18b1e80 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -549,3 +549,58 @@ func TestLateDataIsRolledUp(t *testing.T) { t.Errorf("expected the late point from the minute rollup, got %+v", result) } } + +func TestProcessUsage(t *testing.T) { + st := testStore(t) + ctx := context.Background() + if err := st.AddHost(ctx, "web1", "UTC"); err != nil { + t.Fatal(err) + } + + // php-fpm runs as two workers, postgres is in both lists at once, and + // some snapshots have empty lists + start := time.Now().Add(-time.Hour).Truncate(time.Second) + lists := []monitor.Processes{ + { + CPU: []monitor.Process{ + {Pid: 10, Name: "php-fpm", CPUUsage: 20, MemUsage: 5}, + {Pid: 11, Name: "php-fpm", CPUUsage: 10, MemUsage: 5}, + {Pid: 20, Name: "postgres", CPUUsage: 5, MemUsage: 30}, + }, + Memory: []monitor.Process{ + {Pid: 20, Name: "postgres", CPUUsage: 5, MemUsage: 30}, + {Pid: 10, Name: "php-fpm", CPUUsage: 20, MemUsage: 5}, + }, + }, + {CPU: []monitor.Process{{Pid: 10, Name: "php-fpm", CPUUsage: 40, MemUsage: 5}}}, + {Memory: []monitor.Process{{Pid: 20, Name: "postgres", CPUUsage: 1, MemUsage: 32}}}, + {}, + } + for i, processes := range lists { + snapshot := testSnapshot("web1", start.Add(time.Duration(i)*15*time.Second)) + snapshot.Processes = processes + if err := st.SaveSnapshot(ctx, snapshot); err != nil { + t.Fatal(err) + } + } + + result, err := st.ProcessUsage(ctx, "web1", start, start.Add(time.Minute)) + if err != nil { + t.Fatal(err) + } + want := []ProcessUsage{ + {Name: "php-fpm", CPUAvg: 17.5, CPUPeak: 40, MemAvg: 3.75, MemPeak: 10, SeenPct: 50}, + {Name: "postgres", CPUAvg: 1.5, CPUPeak: 5, MemAvg: 15.5, MemPeak: 32, SeenPct: 50}, + } + if result.Snapshots != 4 || !result.FirstTime.Equal(start) || fmt.Sprint(result.Processes) != fmt.Sprint(want) { + t.Errorf("got %d snapshots from %v: %+v", result.Snapshots, result.FirstTime, result.Processes) + } + + empty, err := st.ProcessUsage(ctx, "web1", start.Add(-time.Hour), start) + if err != nil || empty.Snapshots != 0 || len(empty.Processes) != 0 { + t.Errorf("expected nothing before the first snapshot, got %+v %v", empty, err) + } + if _, err := st.ProcessUsage(ctx, "web1", start, start); !errors.Is(err, ErrInvalid) { + t.Errorf("expected ErrInvalid for an empty range, got %v", err) + } +} From 83e5a44a39efc383a23e9f13d2a662c3f2017541 Mon Sep 17 00:00:00 2001 From: Dhamith Hewamullage Date: Tue, 29 Sep 2026 19:25:34 +0530 Subject: [PATCH 05/10] mark the picked time on charts --- client/web/src/components/ChartCard.svelte | 6 +++-- client/web/src/components/TimeChart.svelte | 27 +++++++++++++++++++++- client/web/src/pages/Host.svelte | 1 + 3 files changed, 31 insertions(+), 3 deletions(-) diff --git a/client/web/src/components/ChartCard.svelte b/client/web/src/components/ChartCard.svelte index f3f0928..f157862 100644 --- a/client/web/src/components/ChartCard.svelte +++ b/client/web/src/components/ChartCard.svelte @@ -19,13 +19,15 @@ note?: string; onzoom?: (from: number, to: number) => void; onpick?: (time: number) => void; + // a picked time to mark on the chart, 0 for none + marker?: number; // replaces the line chart, like the per core heatmap body?: Snippet; // rows for the table view when body replaces the chart table?: Snippet; } - let { title, unit, data, loading = false, error = '', yMax, syncKey, from, to, area = false, note, onzoom, onpick, body, table }: Props = $props(); + let { title, unit, data, loading = false, error = '', yMax, syncKey, from, to, area = false, note, onzoom, onpick, marker = 0, body, table }: Props = $props(); let showTable = $state(false); const colors = ['--series-1', '--series-2', '--series-3', '--series-4', '--series-5', '--series-6', '--series-7', '--series-8']; @@ -98,7 +100,7 @@ {:else if body} {@render body()} {:else if data} - + {/if} {#if data && data.hidden.length > 0} diff --git a/client/web/src/components/TimeChart.svelte b/client/web/src/components/TimeChart.svelte index f450c7e..4864b4c 100644 --- a/client/web/src/components/TimeChart.svelte +++ b/client/web/src/components/TimeChart.svelte @@ -21,9 +21,11 @@ height?: number; onzoom?: (from: number, to: number) => void; onpick?: (time: number) => void; + // a picked time to mark with a line, 0 for none + marker?: number; } - let { data, unit, yMax, syncKey, from, to, area = false, height = 180, onzoom, onpick }: Props = $props(); + let { data, unit, yMax, syncKey, from, to, area = false, height = 180, onzoom, onpick, marker = 0 }: Props = $props(); let wrapper: HTMLDivElement; // uPlot owns this element, Svelte owns the tooltip next to it @@ -94,6 +96,7 @@ })), ], hooks: { + draw: [drawMarker], setCursor: [(u) => updateTooltip(u, colors)], setSelect: [ (u) => { @@ -122,6 +125,23 @@ u.over.style.cursor = onpick ? 'crosshair' : 'default'; } + // a dashed line at the picked time, the moment the process table shows + function drawMarker(u: uPlot) { + if (!marker) return; + const x = Math.round(u.valToPos(marker, 'x', true)); + if (x < u.bbox.left || x > u.bbox.left + u.bbox.width) return; + const ctx = u.ctx; + ctx.save(); + ctx.strokeStyle = cssVar('--accent'); + ctx.lineWidth = uPlot.pxRatio; + ctx.setLineDash([4 * uPlot.pxRatio, 3 * uPlot.pxRatio]); + ctx.beginPath(); + ctx.moveTo(x, u.bbox.top); + ctx.lineTo(x, u.bbox.top + u.bbox.height); + ctx.stroke(); + ctx.restore(); + } + function updateTooltip(u: uPlot, colors: string[]) { const idx = u.cursor.idx; if (!pointerInside || idx == null || u.cursor.left == null || u.cursor.left < 0) { @@ -163,6 +183,11 @@ } }); + $effect(() => { + void marker; + plot?.redraw(false); + }); + onMount(() => { resizeObserver = new ResizeObserver(() => plot?.setSize({ width: wrapper.clientWidth, height })); resizeObserver.observe(wrapper); diff --git a/client/web/src/pages/Host.svelte b/client/web/src/pages/Host.svelte index 029519a..2d38aac 100644 --- a/client/web/src/pages/Host.svelte +++ b/client/web/src/pages/Host.svelte @@ -177,6 +177,7 @@ to={range.to} onzoom={zoom} onpick={pickTime} + marker={processesAt} /> {/each} {#if section.title === 'CPU' && showCores && cores.length === 0} From 6cdce49e72569f49d320ae12cc74417022adfd0d Mon Sep 17 00:00:00 2001 From: Dhamith Hewamullage Date: Tue, 29 Sep 2026 19:28:53 +0530 Subject: [PATCH 06/10] serve prometheus metrics from the client --- client/internal/server/metrics.go | 234 +++++++++++++++++++++++++ client/internal/server/metrics_test.go | 92 ++++++++++ client/internal/server/server.go | 1 + internal/api/api.go | 19 ++ internal/store/query.go | 33 +++- internal/store/store_test.go | 30 ++++ 6 files changed, 408 insertions(+), 1 deletion(-) create mode 100644 client/internal/server/metrics.go create mode 100644 client/internal/server/metrics_test.go diff --git a/client/internal/server/metrics.go b/client/internal/server/metrics.go new file mode 100644 index 0000000..3dc4ec4 --- /dev/null +++ b/client/internal/server/metrics.go @@ -0,0 +1,234 @@ +package server + +import ( + "encoding/json" + "fmt" + "io" + "net/http" + "strconv" + "strings" + + "github.com/dhamith93/SyMon/internal/api" + "github.com/dhamith93/SyMon/internal/logger" + "github.com/dhamith93/SyMon/internal/monitor" +) + +// getMetrics serves every host's latest values in the Prometheus text +// format, so Prometheus can scrape the dashboard +func (s *server) getMetrics(w http.ResponseWriter, r *http.Request) { + response, err := s.collector.Snapshots(r.Context(), &api.Void{}) + if err != nil { + writeGRPCError(w, "snapshots", err) + return + } + + metrics := newMetricsWriter() + for _, host := range response.Hosts { + metrics.gauge("symon_up", "1 when the host reported within the last minute.", boolValue(host.Up), "host", host.Host) + if host.LastSeen > 0 { + metrics.gauge("symon_last_seen_timestamp_seconds", "When the host was last heard from.", float64(host.LastSeen), "host", host.Host) + } + // a host that stopped reporting only has old values, which would + // look current to Prometheus + if !host.Up || host.SnapshotJson == "" { + continue + } + var data monitor.MonitorData + if err := json.Unmarshal([]byte(host.SnapshotJson), &data); err != nil { + logger.Log("error", "cannot read the snapshot of "+host.Host+": "+err.Error()) + continue + } + addHostMetrics(metrics, host.Host, &data) + } + + w.Header().Set("Content-Type", "text/plain; version=0.0.4; charset=utf-8") + if err := metrics.writeTo(w); err != nil { + logger.Log("error", "cannot write metrics: "+err.Error()) + } +} + +const mib = 1024 * 1024 + +func addHostMetrics(m *metricsWriter, host string, data *monitor.MonitorData) { + m.gauge("symon_uptime_seconds", "How long the host has been running.", data.System.UpTimeSeconds, "host", host) + // LoadAvg keeps its old name, it is the CPU usage + m.gauge("symon_cpu_usage_percent", "CPU usage.", float64(data.ProcUsage.LoadAvg), "host", host) + m.gauge("symon_load1", "Load average over 1 minute.", data.ProcUsage.Load1, "host", host) + m.gauge("symon_load5", "Load average over 5 minutes.", data.ProcUsage.Load5, "host", host) + m.gauge("symon_load15", "Load average over 15 minutes.", data.ProcUsage.Load15, "host", host) + + // the agent sends memory and swap in MiB + m.gauge("symon_memory_used_percent", "Memory in use.", data.Memory.PercentageUsed, "host", host) + m.gauge("symon_memory_total_bytes", "Total memory.", float64(data.Memory.Total)*mib, "host", host) + m.gauge("symon_memory_available_bytes", "Memory available to new programs.", float64(data.Memory.Available)*mib, "host", host) + m.gauge("symon_swap_used_percent", "Swap in use.", data.Swap.PercentageUsed, "host", host) + m.gauge("symon_swap_total_bytes", "Total swap.", float64(data.Swap.Total)*mib, "host", host) + + for _, disk := range data.Disk { + labels := []string{"host", host, "device", disk.FileSystem, "mount", disk.MountedOn} + if pct, ok := parsePercent(disk.Usage.Usage); ok { + m.gauge("symon_disk_used_percent", "Disk space in use, as df shows it.", pct, labels...) + } + m.gauge("symon_disk_size_bytes", "Disk size.", float64(disk.Usage.Size), labels...) + m.gauge("symon_disk_used_bytes", "Disk space in use.", float64(disk.Usage.Used), labels...) + if pct, ok := parsePercent(disk.Inodes.Usage); ok { + m.gauge("symon_disk_inodes_used_percent", "Inodes in use.", pct, labels...) + } + } + for _, diskIO := range data.DiskIO { + // named like the disks above + labels := []string{"host", host, "device", "/dev/" + diskIO.Device} + m.gauge("symon_disk_read_bytes_per_second", "Disk reads.", diskIO.ReadBytesPerSec, labels...) + m.gauge("symon_disk_write_bytes_per_second", "Disk writes.", diskIO.WriteBytesPerSec, labels...) + m.gauge("symon_disk_util_percent", "Time the disk was busy.", diskIO.UtilPercent, labels...) + } + for _, network := range data.Networks { + labels := []string{"host", host, "iface", network.Interface} + m.counter("symon_network_receive_bytes_total", "Bytes received since boot.", float64(network.Usage.RxBytes), labels...) + m.counter("symon_network_transmit_bytes_total", "Bytes sent since boot.", float64(network.Usage.TxBytes), labels...) + } + + if tcp := data.TCPStates; tcp != nil { + states := []struct { + name string + count int + }{ + {"established", tcp.Established}, {"syn_sent", tcp.SynSent}, {"syn_recv", tcp.SynRecv}, + {"fin_wait1", tcp.FinWait1}, {"fin_wait2", tcp.FinWait2}, {"time_wait", tcp.TimeWait}, + {"close", tcp.Close}, {"close_wait", tcp.CloseWait}, {"last_ack", tcp.LastAck}, + {"listen", tcp.Listen}, {"closing", tcp.Closing}, + } + for _, state := range states { + m.gauge("symon_tcp_connections", "TCP connections by state.", float64(state.count), "host", host, "state", state.name) + } + } + + if p := data.Pressure; p != nil { + resources := []struct { + name string + pressure monitor.ResourcePressure + }{{"cpu", p.CPU}, {"memory", p.Memory}, {"io", p.IO}} + for _, r := range resources { + m.gauge("symon_pressure_some_percent", "Share of the last 10 seconds some tasks waited on the resource.", r.pressure.Some.Avg10, "host", host, "resource", r.name) + if r.pressure.FullAvailable { + m.gauge("symon_pressure_full_percent", "Share of the last 10 seconds all tasks waited on the resource.", r.pressure.Full.Avg10, "host", host, "resource", r.name) + } + } + } + + for _, temp := range data.Temperatures { + sensor := temp.Name + if temp.Label != "" { + sensor += "/" + temp.Label + } + m.gauge("symon_temperature_celsius", "Sensor temperature.", temp.Celsius, "host", host, "sensor", sensor) + } + for _, service := range data.Services { + m.gauge("symon_service_up", "1 when the service is running.", boolValue(service.Running), "host", host, "service", service.Name) + } + + for _, container := range data.Containers { + name := container.Name + if name == "" { + name = container.ShortID + } + labels := []string{"host", host, "container", name, "project", container.ComposeProject} + m.gauge("symon_container_cpu_percent", "Container CPU usage, as a share of the host.", container.CPU.PercentOfHost, labels...) + m.gauge("symon_container_memory_bytes", "Container memory in use.", container.Memory.Used, labels...) + if rates := container.Rates; rates != nil { + // no traffic rates for containers on the host network + if rates.RxBytesPerSec != nil { + m.gauge("symon_container_receive_bytes_per_second", "Container network traffic received.", *rates.RxBytesPerSec, labels...) + } + if rates.TxBytesPerSec != nil { + m.gauge("symon_container_transmit_bytes_per_second", "Container network traffic sent.", *rates.TxBytesPerSec, labels...) + } + m.gauge("symon_container_read_bytes_per_second", "Container disk reads.", rates.ReadBytesPerSec, labels...) + m.gauge("symon_container_write_bytes_per_second", "Container disk writes.", rates.WriteBytesPerSec, labels...) + } + } +} + +// parsePercent turns "40%" into 40 +func parsePercent(value string) (float64, bool) { + pct, err := strconv.ParseFloat(strings.TrimSuffix(strings.TrimSpace(value), "%"), 64) + return pct, err == nil +} + +func boolValue(b bool) float64 { + if b { + return 1 + } + return 0 +} + +// metricsWriter builds the Prometheus text format. Each metric's HELP, +// TYPE and samples have to be written together, so samples are grouped by +// metric and written at the end. +type metricsWriter struct { + families []*metricFamily + byName map[string]*metricFamily +} + +type metricFamily struct { + name string + kind string + help string + samples []string + // labels already written, a repeated set would make the output invalid + seen map[string]bool +} + +func newMetricsWriter() *metricsWriter { + return &metricsWriter{byName: map[string]*metricFamily{}} +} + +func (w *metricsWriter) gauge(name string, help string, value float64, labels ...string) { + w.add(name, "gauge", help, value, labels) +} + +func (w *metricsWriter) counter(name string, help string, value float64, labels ...string) { + w.add(name, "counter", help, value, labels) +} + +var labelEscaper = strings.NewReplacer(`\`, `\\`, `"`, `\"`, "\n", `\n`) + +// add records one sample. labels are name and value pairs. +func (w *metricsWriter) add(name string, kind string, help string, value float64, labels []string) { + family, ok := w.byName[name] + if !ok { + family = &metricFamily{name: name, kind: kind, help: help, seen: map[string]bool{}} + w.byName[name] = family + w.families = append(w.families, family) + } + + var set strings.Builder + for i := 0; i+1 < len(labels); i += 2 { + if i > 0 { + set.WriteByte(',') + } + fmt.Fprintf(&set, `%s="%s"`, labels[i], labelEscaper.Replace(labels[i+1])) + } + if family.seen[set.String()] { + return + } + family.seen[set.String()] = true + + sample := name + if set.Len() > 0 { + sample += "{" + set.String() + "}" + } + family.samples = append(family.samples, sample+" "+strconv.FormatFloat(value, 'g', -1, 64)+"\n") +} + +func (w *metricsWriter) writeTo(out io.Writer) error { + var b strings.Builder + for _, family := range w.families { + fmt.Fprintf(&b, "# HELP %s %s\n# TYPE %s %s\n", family.name, family.help, family.name, family.kind) + for _, sample := range family.samples { + b.WriteString(sample) + } + } + _, err := io.WriteString(out, b.String()) + return err +} diff --git a/client/internal/server/metrics_test.go b/client/internal/server/metrics_test.go new file mode 100644 index 0000000..c953e76 --- /dev/null +++ b/client/internal/server/metrics_test.go @@ -0,0 +1,92 @@ +package server + +import ( + "context" + "encoding/json" + "net/http/httptest" + "strings" + "testing" + + "github.com/dhamith93/SyMon/internal/api" + "github.com/dhamith93/SyMon/internal/monitor" +) + +// Snapshots has web1 reporting, db1 gone quiet and new1 not sent anything yet +func (f *fakeCollector) Snapshots(ctx context.Context, in *api.Void) (*api.SnapshotList, error) { + snapshot, err := json.Marshal(monitor.MonitorData{ + System: monitor.System{UpTimeSeconds: 3600}, + ProcUsage: monitor.CPU{LoadAvg: 37, Load1: 0.5}, + Memory: monitor.Memory{PercentageUsed: 42.5, Total: 16000, Available: 9000}, + Disk: []monitor.Disk{ + {FileSystem: "/dev/sda1", MountedOn: "/", Usage: monitor.DiskUsage{Size: 1000, Used: 400, Usage: "40%"}, Inodes: monitor.InodeUsage{Usage: "10%"}}, + }, + Networks: []monitor.Network{{Interface: "eth0", Usage: monitor.NetworkUsage{RxBytes: 1000, TxBytes: 2000}}}, + Services: []monitor.Service{{Name: "nginx", Running: true}}, + Containers: []monitor.Container{ + {ShortID: "aaa111", Name: "web", ComposeProject: "shop", CPU: monitor.ContainerCPU{PercentOfHost: 12}, + Rates: &monitor.ContainerRates{RxBytesPerSec: floatPtr(2048), ReadBytesPerSec: 10}}, + }, + }) + if err != nil { + return nil, err + } + return &api.SnapshotList{Hosts: []*api.HostSnapshot{ + {Host: "web1", Up: true, LastSeen: 1700000000, Time: 1700000000, SnapshotJson: string(snapshot)}, + {Host: "db1", LastSeen: 1690000000, Time: 1690000000, SnapshotJson: string(snapshot)}, + {Host: "new1", Up: true, LastSeen: 1700000000}, + }}, nil +} + +func TestMetrics(t *testing.T) { + s, _ := newTestServer(t, nil) + rec := httptest.NewRecorder() + s.routes().ServeHTTP(rec, httptest.NewRequest("GET", "/metrics", nil)) + body := rec.Body.String() + + if rec.Code != 200 || rec.Header().Get("Content-Type") != "text/plain; version=0.0.4; charset=utf-8" { + t.Fatalf("unexpected response %d %q: %s", rec.Code, rec.Header().Get("Content-Type"), body) + } + for _, line := range []string{ + "# TYPE symon_up gauge\n" + `symon_up{host="web1"} 1` + "\n" + `symon_up{host="db1"} 0` + "\n" + `symon_up{host="new1"} 1`, + `symon_last_seen_timestamp_seconds{host="db1"} 1.69e+09`, + `symon_cpu_usage_percent{host="web1"} 37`, + `symon_memory_total_bytes{host="web1"} 1.6777216e+10`, + `symon_disk_used_percent{host="web1",device="/dev/sda1",mount="/"} 40`, + "# TYPE symon_network_receive_bytes_total counter\n" + `symon_network_receive_bytes_total{host="web1",iface="eth0"} 1000`, + `symon_service_up{host="web1",service="nginx"} 1`, + `symon_container_receive_bytes_per_second{host="web1",container="web",project="shop"} 2048`, + } { + if !strings.Contains(body, line+"\n") { + t.Errorf("expected %q in:\n%s", line, body) + } + } + // db1 stopped reporting, so its old values are left out + if strings.Contains(body, `{host="db1",`) || strings.Contains(body, `symon_cpu_usage_percent{host="db1"}`) { + t.Errorf("expected only up and last seen for db1:\n%s", body) + } + // web1 has no transmit rate, as if it were on the host network + if strings.Contains(body, "symon_container_transmit_bytes_per_second") { + t.Errorf("expected no transmit rate:\n%s", body) + } +} + +func TestMetricsWriter(t *testing.T) { + w := newMetricsWriter() + w.gauge("a", "First.", 1, "name", `quote " backslash \ newline`+"\n") + w.counter("b_total", "Second.", 2) + w.gauge("a", "First.", 3, "name", "other") + // a repeated label set is dropped, Prometheus rejects the whole scrape otherwise + w.gauge("a", "First.", 4, "name", "other") + + var out strings.Builder + if err := w.writeTo(&out); err != nil { + t.Fatal(err) + } + want := "# HELP a First.\n# TYPE a gauge\n" + + `a{name="quote \" backslash \\ newline\n"} 1` + "\n" + + `a{name="other"} 3` + "\n" + + "# HELP b_total Second.\n# TYPE b_total counter\nb_total 2\n" + if out.String() != want { + t.Errorf("got:\n%s\nwant:\n%s", out.String(), want) + } +} diff --git a/client/internal/server/server.go b/client/internal/server/server.go index 200bb57..184fe5d 100644 --- a/client/internal/server/server.go +++ b/client/internal/server/server.go @@ -74,6 +74,7 @@ func (s *server) routes() http.Handler { mux.HandleFunc("/api/", func(w http.ResponseWriter, r *http.Request) { writeError(w, http.StatusNotFound, "no such endpoint") }) + mux.HandleFunc("GET /metrics", s.getMetrics) mux.HandleFunc("GET /install.sh", s.getInstallScript) mux.HandleFunc("GET /downloads/{file}", s.getDownload) mux.Handle("/", s.app()) diff --git a/internal/api/api.go b/internal/api/api.go index 3430fa2..8b3cf70 100644 --- a/internal/api/api.go +++ b/internal/api/api.go @@ -221,6 +221,25 @@ func (s *Server) Snapshot(ctx context.Context, in *HostRequest) (*HostSnapshot, }, nil } +func (s *Server) Snapshots(ctx context.Context, in *Void) (*SnapshotList, error) { + hosts, err := s.Store.LatestSnapshots(ctx) + if err != nil { + return nil, toStatus(err) + } + now := time.Now() + list := &SnapshotList{} + for _, latest := range hosts { + list.Hosts = append(list.Hosts, &HostSnapshot{ + Host: latest.Host, + Time: unix(latest.Time), + SnapshotJson: string(latest.Snapshot), + LastSeen: unix(latest.LastSeen), + Up: isUp(latest.LastSeen, now), + }) + } + return list, nil +} + func (s *Server) QuerySeries(ctx context.Context, in *SeriesRequest) (*SeriesResponse, error) { result, err := s.Store.QuerySeries(ctx, store.SeriesQuery{ Host: in.Host, diff --git a/internal/store/query.go b/internal/store/query.go index 7a5bc80..5c4decf 100644 --- a/internal/store/query.go +++ b/internal/store/query.go @@ -92,6 +92,8 @@ func fillSummary(summary *HostSummary, data *monitor.MonitorData) { } type LatestSnapshot struct { + Host string + // Time is zero, and Snapshot empty, when the host never sent one Time time.Time // Snapshot is the MonitorData as the agent sent it Snapshot []byte @@ -100,7 +102,7 @@ type LatestSnapshot struct { // LatestSnapshot returns a host's newest snapshot and when it was last heard from func (s *Store) LatestSnapshot(ctx context.Context, host string) (LatestSnapshot, error) { - var latest LatestSnapshot + latest := LatestSnapshot{Host: host} var lastSeen *time.Time err := s.pool.QueryRow(ctx, ` SELECT l.time, l.snapshot, h.last_seen FROM host_latest l @@ -115,6 +117,35 @@ func (s *Store) LatestSnapshot(ctx context.Context, host string) (LatestSnapshot return latest, err } +// LatestSnapshots returns every host with its newest snapshot +func (s *Store) LatestSnapshots(ctx context.Context) ([]LatestSnapshot, error) { + rows, err := s.pool.Query(ctx, ` + SELECT h.name, h.last_seen, l.time, l.snapshot FROM hosts h + LEFT JOIN host_latest l ON l.host_id = h.id + ORDER BY h.name`) + if err != nil { + return nil, err + } + defer rows.Close() + + hosts := []LatestSnapshot{} + for rows.Next() { + var latest LatestSnapshot + var lastSeen, snapshotTime *time.Time + if err := rows.Scan(&latest.Host, &lastSeen, &snapshotTime, &latest.Snapshot); err != nil { + return nil, err + } + if lastSeen != nil { + latest.LastSeen = *lastSeen + } + if snapshotTime != nil { + latest.Time = *snapshotTime + } + hosts = append(hosts, latest) + } + return hosts, rows.Err() +} + // Processes returns the top process lists from the newest snapshot at or // before the given time func (s *Store) Processes(ctx context.Context, host string, at time.Time) (time.Time, []byte, error) { diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 18b1e80..be1e5ea 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -7,6 +7,7 @@ import ( "fmt" "os" "strconv" + "strings" "testing" "time" @@ -604,3 +605,32 @@ func TestProcessUsage(t *testing.T) { t.Errorf("expected ErrInvalid for an empty range, got %v", err) } } + +func TestLatestSnapshots(t *testing.T) { + st := testStore(t) + ctx := context.Background() + for _, host := range []string{"web1", "new1"} { + if err := st.AddHost(ctx, host, "UTC"); err != nil { + t.Fatal(err) + } + } + at := time.Now().Add(-time.Minute).Truncate(time.Second) + if err := st.SaveSnapshot(ctx, testSnapshot("web1", at)); err != nil { + t.Fatal(err) + } + + hosts, err := st.LatestSnapshots(ctx) + if err != nil { + t.Fatal(err) + } + if len(hosts) != 2 { + t.Fatalf("expected two hosts, got %+v", hosts) + } + // new1 has not sent anything yet + if hosts[0].Host != "new1" || !hosts[0].Time.IsZero() || len(hosts[0].Snapshot) != 0 { + t.Errorf("expected new1 without a snapshot, got %+v", hosts[0]) + } + if hosts[1].Host != "web1" || !hosts[1].Time.Equal(at) || !strings.Contains(string(hosts[1].Snapshot), "Debian 13") { + t.Errorf("expected web1's snapshot, got %+v", hosts[1]) + } +} From 5837c1657b48d095b1f0c76da398b8aeb90c410f Mon Sep 17 00:00:00 2001 From: Dhamith Hewamullage Date: Tue, 29 Sep 2026 19:30:00 +0530 Subject: [PATCH 07/10] document forecasts, process usage and /metrics --- README.MD | 16 ++++++++++++---- docs/install.md | 18 ++++++++++++++++++ 2 files changed, 30 insertions(+), 4 deletions(-) diff --git a/README.MD b/README.MD index 84a26e0..da39602 100644 --- a/README.MD +++ b/README.MD @@ -11,8 +11,9 @@ SyMon is a self-hosted monitoring tool for Linux servers, home labs and Raspberr **Hosts** - CPU, overall and per core, load average, memory and swap - Disk space, disk IO and busy time, network traffic, TCP connections +- When each disk will be full, at the rate it grew over the last week - Pressure stall information and temperature sensors -- The top processes by CPU and memory, at any point in time +- The top processes by CPU and memory at any point in time, and the programs that used the most over any range - Whether chosen systemd services are running **Containers** @@ -23,18 +24,19 @@ SyMon is a self-hosted monitoring tool for Linux servers, home labs and Raspberr - Send any number from a script or cron job and get a chart for it **Alerts** -- Rules for CPU, memory, swap, disks, services, custom metrics, silent hosts and HTTP endpoints +- Rules for CPU, memory, swap, disks, disks filling up, services, custom metrics, silent hosts and HTTP endpoints - Warning and critical levels, shown on the dashboard and sent by email, Slack or PagerDuty **Dashboard** - Every host at a glance, and a page per host with charts from 15 minutes to 30 days, or any custom range -- Drag across a chart to zoom in, and switch any chart to a table +- Drag across a chart to zoom in, click a point to see the processes running then, and switch any chart to a table - Follows the system light or dark theme **Running it** - Add a host with one command, with a single-use token - Raw data kept for 7 days, 1 minute averages for 30 days and 1 hour averages for a year, all adjustable - Each host has its own key, and components can talk over TLS +- A Prometheus endpoint, for Grafana or a Prometheus you already run ## Screenshots @@ -96,7 +98,7 @@ Components talk over gRPC, so other tools can read from or push into them. See t The Client exposes a JSON API under `/api/v1`. Times are unix seconds. Errors return a JSON body `{"error": "..."}` with a 4xx or 5xx status. * `GET /api/v1/fleet` - * Every host with its status, latest usage, number of running containers and number of open alerts + * Every host with its status, latest usage, number of running containers, number of open alerts, and `diskFullDays`, the days until its first disk is full (null when none is filling up) * `GET /api/v1/hosts/{host}` * The host's latest snapshot as the agent sent it, with `up` and `lastSeen` * `GET /api/v1/hosts/{host}/series?metric=cpu&from=&to=` @@ -104,7 +106,13 @@ The Client exposes a JSON API under `/api/v1`. Times are unix seconds. Errors re * Metrics: `cpu`, `cpu_core`, `load1`, `load5`, `load15`, `memory`, `memory_used`, `swap`, `swap_used`, `psi_cpu`, `psi_memory`, `psi_memory_full`, `psi_io`, `psi_io_full`, `tcp_established`, `tcp_time_wait`, `tcp_close_wait`, `tcp_listen`, `tcp_total`, `disk_used`, `disk_inodes`, `disk_read`, `disk_write`, `disk_util`, `net_rx`, `net_tx`, `temperature`, `custom`, `container_cpu`, `container_memory`, `container_rx`, `container_tx`, `container_io_read`, `container_io_write` * `GET /api/v1/hosts/{host}/processes?at=` * Top processes by CPU and by memory at or before `at`, or the latest +* `GET /api/v1/hosts/{host}/process-usage?from=&to=` + * The programs that used the most CPU and memory over the range, with their average, peak and how often they were among the top processes. Processes with the same name are added up +* `GET /api/v1/hosts/{host}/disk-forecasts` + * Each disk's growth per day over the last week, and `daysToFull`, or null when the disk is not filling up * `GET /api/v1/hosts/{host}/custom-metrics` * Names of the host's custom metrics * `GET /api/v1/alerts?host=&open=1&from=&to=` * Alerts, newest first. `open=1` leaves out resolved ones + +The Client also serves `GET /metrics`, every host's latest values in the Prometheus text format. See [Prometheus and Grafana](docs/install.md#prometheus-and-grafana). diff --git a/docs/install.md b/docs/install.md index ae1a993..7c21923 100644 --- a/docs/install.md +++ b/docs/install.md @@ -6,6 +6,7 @@ This guide sets up SyMon on one Linux server and adds hosts to it. It covers upg - [Install the server](#install-the-server) - [Add hosts](#add-hosts) - [Alerts](#alerts) +- [Prometheus and Grafana](#prometheus-and-grafana) - [Upgrades](#upgrades) - [Backups](#backups) - [Uninstall](#uninstall) @@ -212,6 +213,7 @@ A rule looks like this: | `memory` | memory used, % | | | `swap` | swap used, % | | | `disks` | disk space used, % | `Disk`, the device | +| `disk_forecast` | days until the disk is full, at its growth over the last week | `Disk`, the device | | `services` | a service from the service list. `Op` `inactive` alerts when it stops, `active` when it runs | `Service`, the name from the service list | | `ping` | host silent for longer than `TriggerIntveral` seconds | | | `endpoint` | an HTTP check from the collector | `Endpoint`, `Method`, `ExpectedHTTPCode`, `POSTBody`, `POSTContentType` | @@ -219,6 +221,8 @@ A rule looks like this: `Op` is one of `>`, `<`, `>=`, `<=`, `==` or `!=`. A value has to stay past a threshold for `TriggerIntveral` seconds before the alert opens, and back to normal for as long before it resolves. Endpoint checks need `SYMON_ENABLE_ENDPOINT_MONITORING=true` on the collector. +`disk_forecast` rules use `Op` `<`, for example a warning under 14 days and critical under 3. A forecast needs a day of history and steady growth, so a disk that fills and empties, like one with rotating logs, gets none. A disk that is not filling up counts as 365 days. The dashboard shows the forecast in the host's disk table, and on the hosts page when a disk fills up within 30 days. + The collector reads the rules when it starts, so restart it after editing them. ### The alert processor @@ -243,6 +247,20 @@ sudo systemctl daemon-reload sudo systemctl enable --now symon_alertprocessor ``` +## Prometheus and Grafana + +The dashboard serves every host's latest values at `/metrics` in the Prometheus format, so Grafana or an existing Prometheus can use them. Add it to `prometheus.yml`: + +```yaml +scrape_configs: + - job_name: symon + scrape_interval: 15s + static_configs: + - targets: ["symon.example.lan:8080"] +``` + +Every value has a `host` label. Disks, interfaces, sensors, services and containers have their own labels too. `symon_up` is 0 for a host that stopped reporting, and its other values are left out until it reports again. Like the rest of the dashboard, `/metrics` has no login, so keep it behind the same reverse proxy or firewall. + ## Upgrades Upgrade the server first, then the hosts. From f50dd1a35f9cd9a672cc1d85b83728004a1bda47 Mon Sep 17 00:00:00 2001 From: Dhamith Hewamullage Date: Tue, 29 Sep 2026 21:04:56 +0530 Subject: [PATCH 08/10] limit process usage to a day --- README.MD | 2 +- .../src/components/BusiestProcesses.svelte | 10 ++- internal/store/load_test.go | 65 +++++++++++++++++++ internal/store/query.go | 7 ++ internal/store/store_test.go | 3 + 5 files changed, 83 insertions(+), 4 deletions(-) diff --git a/README.MD b/README.MD index da39602..848b296 100644 --- a/README.MD +++ b/README.MD @@ -107,7 +107,7 @@ The Client exposes a JSON API under `/api/v1`. Times are unix seconds. Errors re * `GET /api/v1/hosts/{host}/processes?at=` * Top processes by CPU and by memory at or before `at`, or the latest * `GET /api/v1/hosts/{host}/process-usage?from=&to=` - * The programs that used the most CPU and memory over the range, with their average, peak and how often they were among the top processes. Processes with the same name are added up + * The programs that used the most CPU and memory over the range, 24 hours at most, with their average, peak and how often they were among the top processes. Processes with the same name are added up * `GET /api/v1/hosts/{host}/disk-forecasts` * Each disk's growth per day over the last week, and `daysToFull`, or null when the disk is not filling up * `GET /api/v1/hosts/{host}/custom-metrics` diff --git a/client/web/src/components/BusiestProcesses.svelte b/client/web/src/components/BusiestProcesses.svelte index 925b149..4f53c64 100644 --- a/client/web/src/components/BusiestProcesses.svelte +++ b/client/web/src/components/BusiestProcesses.svelte @@ -6,6 +6,9 @@ let { host, from, to }: { host: string; from: number; to: number } = $props(); const limit = 15; + // the collector answers for up to a day, so a longer range shows its last day + const maxSpan = 24 * 3600; + const start = $derived(Math.max(from, to - maxSpan)); let view = $state<'CPU' | 'Memory'>('CPU'); let snapshots = $state(0); @@ -18,7 +21,7 @@ let requestId = 0; $effect(() => { - const request = { host, from, to, id: ++requestId }; + const request = { host, from: start, to, id: ++requestId }; loading = true; api .processUsage(request.host, request.from, request.to) @@ -55,8 +58,9 @@

{#if snapshots > 0} - From {snapshots.toLocaleString()} snapshots since {formatDateTime(firstTime)}. Each keeps only the top 10 processes, so - averages are a lower bound. + {start > from ? 'The last 24 hours of this range, from' : 'From'} + {snapshots.toLocaleString()} snapshots since {formatDateTime(firstTime)}. Each keeps only the top 10 processes, so averages + are a lower bound. {:else if !error && !loading} No process lists recorded in this range. {/if} diff --git a/internal/store/load_test.go b/internal/store/load_test.go index 76ab5e8..9598756 100644 --- a/internal/store/load_test.go +++ b/internal/store/load_test.go @@ -11,6 +11,7 @@ import ( "testing" "time" + "github.com/dhamith93/SyMon/internal/monitor" "github.com/jackc/pgx/v5" ) @@ -127,3 +128,67 @@ func TestLoadQuery30Days(t *testing.T) { t.Errorf("query took %v, over 100ms", best) } } + +// a host sending its top processes every 15s for 7 days, the raw +// retention, then the busiest processes over the last day and over a day +// in compressed chunks +func TestLoadProcessUsage7Days(t *testing.T) { + st := loadStore(t) + ctx := context.Background() + if err := st.AddHost(ctx, "web1", "UTC"); err != nil { + t.Fatal(err) + } + hostID, err := st.hostID(ctx, "web1") + if err != nil { + t.Fatal(err) + } + + // 30 programs, top 10 by CPU and by memory in each snapshot, some in both + process := func() monitor.Process { + n := rand.Intn(30) + return monitor.Process{Pid: 1000 + n*10 + rand.Intn(3), Name: fmt.Sprintf("program%02d", n), ExecPath: fmt.Sprintf("/usr/bin/program%02d", n), + User: "root", CPUUsage: rand.Float32() * 50, MemUsage: rand.Float32() * 10, Threads: 4} + } + end := time.Now().Truncate(time.Second) + start := end.Add(-7 * 24 * time.Hour) + var rows [][]any + for at := start; at.Before(end); at = at.Add(15 * time.Second) { + var processes monitor.Processes + for i := 0; i < 10; i++ { + processes.CPU = append(processes.CPU, process()) + processes.Memory = append(processes.Memory, process()) + } + rows = append(rows, []any{at, hostID, processes}) + } + copied, err := st.pool.CopyFrom(ctx, pgx.Identifier{"process_snapshots"}, []string{"time", "host_id", "processes"}, pgx.CopyFromRows(rows)) + if err != nil { + t.Fatal(err) + } + // the compression policy would have compressed these by now + var compressed int + err = st.pool.QueryRow(ctx, `SELECT count(compress_chunk(c)) FROM show_chunks('process_snapshots', older_than => now() - INTERVAL '2 days') c`).Scan(&compressed) + if err != nil { + t.Fatal(err) + } + t.Logf("inserted %d snapshots, compressed %d chunks", copied, compressed) + + for _, daysAgo := range []int{0, 5} { + to := end.Add(-time.Duration(daysAgo) * 24 * time.Hour) + var best time.Duration + var result ProcessUsageResult + for i := 0; i < 3; i++ { + started := time.Now() + result, err = st.ProcessUsage(ctx, "web1", to.Add(-processUsageMaxRange), to) + if err != nil { + t.Fatal(err) + } + if took := time.Since(started); best == 0 || took < best { + best = took + } + } + t.Logf("day ending %d days ago: %d snapshots, %d programs, best of 3 %v", daysAgo, result.Snapshots, len(result.Processes), best) + if best > time.Second { + t.Errorf("day ending %d days ago took %v, over 1s", daysAgo, best) + } + } +} diff --git a/internal/store/query.go b/internal/store/query.go index 5c4decf..c0ce6c1 100644 --- a/internal/store/query.go +++ b/internal/store/query.go @@ -189,6 +189,10 @@ type ProcessUsageResult struct { // CPU and memory const processUsageLimit = 15 +// processUsageMaxRange keeps ProcessUsage under a second. A day of +// snapshots takes about 0.3s on the dev VM, a week over 2s. +const processUsageMaxRange = 24 * time.Hour + // ProcessUsage adds up each program's CPU and memory over a range. Snapshots // only keep the top processes, so a program counts as 0 where it was not // among them, and the averages are a lower bound. @@ -196,6 +200,9 @@ func (s *Store) ProcessUsage(ctx context.Context, host string, from time.Time, t if !to.After(from) { return ProcessUsageResult{}, fmt.Errorf("%w: from must be before to", ErrInvalid) } + if to.Sub(from) > processUsageMaxRange { + return ProcessUsageResult{}, fmt.Errorf("%w: pick 24 hours or less", ErrInvalid) + } hostID, err := s.hostID(ctx, host) if err != nil { return ProcessUsageResult{}, err diff --git a/internal/store/store_test.go b/internal/store/store_test.go index be1e5ea..2c373f0 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -604,6 +604,9 @@ func TestProcessUsage(t *testing.T) { if _, err := st.ProcessUsage(ctx, "web1", start, start); !errors.Is(err, ErrInvalid) { t.Errorf("expected ErrInvalid for an empty range, got %v", err) } + if _, err := st.ProcessUsage(ctx, "web1", start.Add(-25*time.Hour), start); !errors.Is(err, ErrInvalid) { + t.Errorf("expected ErrInvalid for a range over a day, got %v", err) + } } func TestLatestSnapshots(t *testing.T) { From f1c96a642e420b9e28ca8f12a70378dbe5f52e4e Mon Sep 17 00:00:00 2001 From: Dhamith Hewamullage Date: Tue, 29 Sep 2026 21:10:59 +0530 Subject: [PATCH 09/10] fit disk forecasts on bytes --- internal/store/forecast.go | 34 ++++++++++++++++----------- internal/store/forecast_test.go | 41 ++++++++++++++++++--------------- 2 files changed, 44 insertions(+), 31 deletions(-) diff --git a/internal/store/forecast.go b/internal/store/forecast.go index 78c09b7..8c28dc1 100644 --- a/internal/store/forecast.go +++ b/internal/store/forecast.go @@ -8,7 +8,8 @@ import ( // A disk is forecast to fill up only when its usage has grown steadily over // the window. Disks that jump up and down, like ones with rotating logs, -// have a low r2 and get no forecast. +// have a low r2 and get no forecast. The fit is on bytes, since the agent +// sends df's percent as a whole number, which hides slow growth. const ( forecastWindow = 7 * 24 * time.Hour forecastMinSamples = 24 @@ -34,16 +35,20 @@ type DiskForecast struct { } // forecastDays returns how many days are left until a disk is full, and -// false when the disk is not filling up or there is too little history to say -func forecastDays(usedPct float64, pctPerDay float64, r2 float64, samples int) (float64, bool) { - if samples < forecastMinSamples || pctPerDay <= 0 || r2 < forecastMinR2 { +// false when the disk is not filling up or there is too little history to +// say. df's percent is used / (used + available), where available leaves +// out the blocks reserved for root, so the space left comes from it rather +// than from the disk size. +func forecastDays(usedPct float64, usedBytes float64, bytesPerDay float64, r2 float64, samples int) (float64, bool) { + if samples < forecastMinSamples || bytesPerDay <= 0 || r2 < forecastMinR2 || usedPct <= 0 { return 0, false } - days := (100 - usedPct) / pctPerDay + left := max(usedBytes*(100-usedPct)/usedPct, 0) + days := left / bytesPerDay if days > forecastHorizonDays { return 0, false } - return max(days, 0), true + return days, true } // DiskForecasts returns a forecast for each of a host's disks @@ -112,10 +117,10 @@ func (s *Store) queryForecasts(ctx context.Context, filter string, args ...any) sql := fmt.Sprintf(` SELECT h.name, d.device, d.mount, last(d.used_pct, d.bucket), - regr_slope(d.used_pct, extract(epoch FROM d.bucket - $1::timestamptz) / 86400), + last(d.used_bytes, d.bucket), regr_slope(d.used_bytes, extract(epoch FROM d.bucket - $1::timestamptz) / 86400), - regr_r2(d.used_pct, extract(epoch FROM d.bucket - $1::timestamptz) / 86400), - count(d.used_pct) + regr_r2(d.used_bytes, extract(epoch FROM d.bucket - $1::timestamptz) / 86400), + count(d.used_bytes) FROM disk_metrics_1h d JOIN hosts h ON h.id = d.host_id WHERE d.bucket >= $1 %s @@ -130,15 +135,18 @@ func (s *Store) queryForecasts(ctx context.Context, filter string, args ...any) forecasts := []DiskForecast{} for rows.Next() { var forecast DiskForecast - var usedPct, pctPerDay, bytesPerDay, r2 *float64 - if err := rows.Scan(&forecast.Host, &forecast.Device, &forecast.Mount, &usedPct, &pctPerDay, &bytesPerDay, &r2, &forecast.Samples); err != nil { + var usedPct, usedBytes, bytesPerDay, r2 *float64 + if err := rows.Scan(&forecast.Host, &forecast.Device, &forecast.Mount, &usedPct, &usedBytes, &bytesPerDay, &r2, &forecast.Samples); err != nil { return nil, err } // nulls come from a single sample or no usable values forecast.UsedPct = valueOrZero(usedPct) - forecast.PctPerDay = valueOrZero(pctPerDay) forecast.BytesPerDay = valueOrZero(bytesPerDay) - if days, ok := forecastDays(forecast.UsedPct, forecast.PctPerDay, valueOrZero(r2), forecast.Samples); ok { + // the growth in df's percent points, at the current size + if used := valueOrZero(usedBytes); used > 0 { + forecast.PctPerDay = forecast.BytesPerDay * forecast.UsedPct / used + } + if days, ok := forecastDays(forecast.UsedPct, valueOrZero(usedBytes), forecast.BytesPerDay, valueOrZero(r2), forecast.Samples); ok { forecast.DaysToFull = &days } forecasts = append(forecasts, forecast) diff --git a/internal/store/forecast_test.go b/internal/store/forecast_test.go index 8364aee..3786455 100644 --- a/internal/store/forecast_test.go +++ b/internal/store/forecast_test.go @@ -11,26 +11,30 @@ import ( ) func TestForecastDays(t *testing.T) { + // 600 GB used at 60% leaves 400 GB, the rest is reserved or already used + const gb = 1e9 tests := []struct { - name string - usedPct float64 - pctPerDay float64 - r2 float64 - samples int - wantDays float64 - wantOK bool + name string + usedPct float64 + usedBytes float64 + bytesPerDay float64 + r2 float64 + samples int + wantDays float64 + wantOK bool }{ - {"steady growth", 60, 2, 0.95, 168, 20, true}, - {"flat", 60, 0, 1, 168, 0, false}, - {"shrinking", 60, -1, 0.9, 168, 0, false}, - {"noisy", 60, 2, 0.3, 168, 0, false}, - {"one day of history", 60, 2, 0.95, 24, 20, true}, - {"too little history", 60, 2, 0.95, 23, 0, false}, - {"beyond the horizon", 10, 0.2, 0.95, 168, 0, false}, - {"already full", 100, 1, 0.95, 168, 0, true}, + {"steady growth", 60, 600 * gb, 20 * gb, 0.95, 168, 20, true}, + {"flat", 60, 600 * gb, 0, 1, 168, 0, false}, + {"shrinking", 60, 600 * gb, -10 * gb, 0.9, 168, 0, false}, + {"noisy", 60, 600 * gb, 20 * gb, 0.3, 168, 0, false}, + {"one day of history", 60, 600 * gb, 20 * gb, 0.95, 24, 20, true}, + {"too little history", 60, 600 * gb, 20 * gb, 0.95, 23, 0, false}, + {"beyond the horizon", 60, 600 * gb, 1 * gb, 0.95, 168, 0, false}, + {"already full", 100, 600 * gb, 1 * gb, 0.95, 168, 0, true}, + {"empty disk", 0, 0, 1 * gb, 0.95, 168, 0, false}, } for _, tt := range tests { - days, ok := forecastDays(tt.usedPct, tt.pctPerDay, tt.r2, tt.samples) + days, ok := forecastDays(tt.usedPct, tt.usedBytes, tt.bytesPerDay, tt.r2, tt.samples) if ok != tt.wantOK || days != tt.wantDays { t.Errorf("%s: got %v %v, want %v %v", tt.name, days, ok, tt.wantDays, tt.wantOK) } @@ -68,14 +72,15 @@ func TestDiskForecasts(t *testing.T) { } // two days of samples: /data grows 2% a day and reaches 60% now, so it - // is full in about 20 days. / stays at 40%. + // is full in about 20 days. / stays at 40%. Percents are whole numbers, + // like the agent sends them. now := time.Now() batch := &pgx.Batch{} insert := `INSERT INTO disk_metrics (time, host_id, device, mount, fstype, size_bytes, used_bytes, used_pct, inodes_used_pct) VALUES ($1, $2, $3, $4, 'ext4', 1e9, $5, $6, 1)` for at := now.Add(-48 * time.Hour); at.Before(now); at = at.Add(10 * time.Minute) { used := 60 + 2*at.Sub(now).Hours()/24 - batch.Queue(insert, at, hostID, "/dev/sdb1", "/data", used*1e7, used) + batch.Queue(insert, at, hostID, "/dev/sdb1", "/data", used*1e7, math.Round(used)) batch.Queue(insert, at, hostID, "/dev/sda1", "/", 40e7, 40.0) } if err := st.sendBatch(ctx, batch); err != nil { From df037f9c5e3f2e74c16752da9de436991c48c814 Mon Sep 17 00:00:00 2001 From: Dhamith Hewamullage Date: Tue, 29 Sep 2026 21:20:24 +0530 Subject: [PATCH 10/10] update installation instructions to include PostgreSQL package --- docs/install.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/install.md b/docs/install.md index 7c21923..3d7b420 100644 --- a/docs/install.md +++ b/docs/install.md @@ -63,7 +63,7 @@ echo "deb https://packagecloud.io/timescale/timescaledb/ubuntu/ $(lsb_release -c wget -qO- https://packagecloud.io/timescale/timescaledb/gpgkey \ | sudo gpg --dearmor -o /etc/apt/trusted.gpg.d/timescaledb.gpg sudo apt update -sudo apt install -y timescaledb-2-postgresql-18 +sudo apt install -y postgresql-18 timescaledb-2-postgresql-18 sudo timescaledb-tune --quiet --yes sudo systemctl restart postgresql ```