// p4changestats summarizes submitted changelists and file revisions by month. package main import ( "bufio" "bytes" "encoding/base64" "encoding/csv" "encoding/json" "errors" "fmt" "io" "net/http" "net/url" "os" "os/exec" "path/filepath" "sort" "strconv" "strings" "time" "github.com/rcowham/kingpin" "github.com/sirupsen/logrus" ) const ( vmagentEnvFile = "/var/vmagent/vmagent.env" vmagentPasswordFile = "/var/vmagent/.vmpassword" description = "Summarize Perforce submitted changes and file revisions by calendar month." defaultPushTimeout = 60 * time.Second ) // Version can be set at build time with -ldflags "-X main.Version=". var Version = "v0.3" type reportRow struct { Month string `json:"month"` Changes int `json:"changes"` AddEdit int `json:"add_edit"` Deleted int `json:"deleted"` Branched int `json:"branched"` Purged int `json:"purged"` } type revisionCounts struct { addEdit int deleted int branched int purged int } type depotInfo struct { name string typev string streamDepth int } type trackingTotals struct { commands int lapse time.Duration userCPU time.Duration systemCPU time.Duration memoryMB float64 readLockHeld time.Duration writeLockHeld time.Duration tableLockHeld map[string]lockHeld currentTable string } type lockHeld struct { read time.Duration write time.Duration } type trackingReport struct { Commands int `json:"commands"` Lapse string `json:"lapse"` AverageLapse string `json:"average_lapse"` UserCPU string `json:"user_cpu"` AverageUserCPU string `json:"average_user_cpu"` SystemCPU string `json:"system_cpu"` AverageSystemCPU string `json:"average_system_cpu"` MemoryMB float64 `json:"memory_mb"` AverageMemoryMB float64 `json:"average_memory_mb"` ReadLockHeld string `json:"read_lock_held"` WriteLockHeld string `json:"write_lock_held"` Tables []tableLockReport `json:"tables"` } type tableLockReport struct { Table string `json:"table"` ReadHeld string `json:"read_lock_held"` WriteHeld string `json:"write_lock_held"` } type remoteReport struct { Version string `json:"version"` Rows []reportRow `json:"rows"` Tracking trackingReport `json:"tracking"` } type p4Runner struct { bin string port string user string client string logger *logrus.Logger tracking trackingTotals streamDirectories map[string][]string } func (p *p4Runner) run(args ...string) ([]byte, error) { commandArgs := p.commandArgs(args...) p.logger.Debugf("Running: %s %s", p.bin, strings.Join(commandArgs, " ")) command := exec.Command(p.bin, commandArgs...) p.tracking.commands++ var stderr bytes.Buffer command.Stderr = &stderr output, err := command.Output() if err != nil { return nil, p.commandError(args, err, stderr.String()) } p.tracking.add(output) p.logger.Debug(p.tracking.String()) return output, nil } func (p *p4Runner) commandArgs(args ...string) []string { commandArgs := make([]string, 0, len(args)+6) if p.port != "" { commandArgs = append(commandArgs, "-p", p.port) } if p.user != "" { commandArgs = append(commandArgs, "-u", p.user) } if p.client != "" { commandArgs = append(commandArgs, "-c", p.client) } commandArgs = append(commandArgs, "-Ztrack", "-ztag") commandArgs = append(commandArgs, args...) return commandArgs } func (p *p4Runner) commandError(args []string, err error, stderr string) error { message := strings.TrimSpace(stderr) if message != "" { return fmt.Errorf("p4 %s: %w: %s", strings.Join(args, " "), err, message) } return fmt.Errorf("p4 %s: %w", strings.Join(args, " "), err) } func (p *p4Runner) streamFileRevisionCounts(filespec string) (revisionCounts, error) { args := []string{"files", "-a", filespec} commandArgs := p.commandArgs(args...) p.logger.Debugf("Running: %s %s", p.bin, strings.Join(commandArgs, " ")) command := exec.Command(p.bin, commandArgs...) p.tracking.commands++ stdout, err := command.StdoutPipe() if err != nil { return revisionCounts{}, err } var stderr bytes.Buffer command.Stderr = &stderr if err := command.Start(); err != nil { return revisionCounts{}, p.commandError(args, err, stderr.String()) } counts := revisionCounts{} scanner := bufio.NewScanner(stdout) scanner.Buffer(make([]byte, 64*1024), 4*1024*1024) for scanner.Scan() { line := scanner.Text() if strings.HasPrefix(line, "... action ") { counts.addAction(strings.TrimPrefix(line, "... action ")) continue } p.tracking.addLine(line) } scanErr := scanner.Err() waitErr := command.Wait() if scanErr != nil { return revisionCounts{}, fmt.Errorf("read p4 files output: %w", scanErr) } if waitErr != nil { return revisionCounts{}, p.commandError(args, waitErr, stderr.String()) } p.logger.Debug(p.tracking.String()) return counts, nil } func (c *revisionCounts) addAction(action string) { switch action { case "add", "edit", "move/add", "integrate": c.addEdit++ case "delete", "move/delete": c.deleted++ case "branch", "import": c.branched++ case "purge": c.purged++ } } func (t *trackingTotals) add(output []byte) { t.currentTable = "" for _, line := range strings.Split(string(output), "\n") { t.addLine(line) } } func (t *trackingTotals) addLine(line string) { if t.tableLockHeld == nil { t.tableLockHeld = make(map[string]lockHeld) } line = strings.TrimSpace(line) if !strings.HasPrefix(line, "--- ") { return } track := strings.TrimSpace(strings.TrimPrefix(line, "--- ")) switch { case strings.HasPrefix(track, "lapse "): t.currentTable = "" t.lapse += parseTrackingDuration(strings.TrimSpace(strings.TrimPrefix(track, "lapse "))) case strings.HasPrefix(track, "usage "): fields := strings.Fields(strings.TrimPrefix(track, "usage ")) if len(fields) > 0 { userCPU, systemCPU, found := strings.Cut(fields[0], "+") if found { t.userCPU += parseCPUTime(userCPU) t.systemCPU += parseCPUTime(systemCPU) } } case strings.HasPrefix(track, "memory cmd/proc "): memory := strings.TrimSpace(strings.TrimPrefix(track, "memory cmd/proc ")) cmdMemory, _, _ := strings.Cut(memory, "/") t.memoryMB += parseMemoryMB(cmdMemory) case strings.HasPrefix(track, "db."): t.currentTable = strings.Fields(track)[0] case strings.HasPrefix(track, "total lock wait+held read/write ") && t.currentTable != "": lockTimes := strings.TrimSpace(strings.TrimPrefix(track, "total lock wait+held read/write ")) readTimes, writeTimes, found := strings.Cut(lockTimes, "/") if !found { return } _, readHeld, readFound := strings.Cut(readTimes, "+") _, writeHeld, writeFound := strings.Cut(writeTimes, "+") if !readFound || !writeFound { return } readHeldDuration := parseTrackingDuration(readHeld) writeHeldDuration := parseTrackingDuration(writeHeld) t.readLockHeld += readHeldDuration t.writeLockHeld += writeHeldDuration held := t.tableLockHeld[t.currentTable] held.read += readHeldDuration held.write += writeHeldDuration t.tableLockHeld[t.currentTable] = held } } func parseTrackingDuration(value string) time.Duration { duration, err := time.ParseDuration(value) if err != nil { return 0 } return duration } func parseCPUTime(value string) time.Duration { value = strings.TrimRight(value, "abcdefghijklmnopqrstuvwxyz") milliseconds, err := strconv.ParseFloat(value, 64) if err != nil { return 0 } return time.Duration(milliseconds * float64(time.Millisecond)) } func parseMemoryMB(value string) float64 { value = strings.ToLower(strings.TrimSpace(value)) for _, unit := range []string{"kb", "mb", "gb"} { if number, found := strings.CutSuffix(value, unit); found { amount, err := strconv.ParseFloat(number, 64) if err != nil { return 0 } switch unit { case "kb": return amount / 1024 case "gb": return amount * 1024 default: return amount } } } return 0 } func formatAverageMilliseconds(duration time.Duration) string { return fmt.Sprintf("%.1fms", float64(duration)/float64(time.Millisecond)) } func (t trackingTotals) String() string { report := t.report() return fmt.Sprintf("P4 tracking totals: commands=%d lapse=%s average_lapse=%s user_cpu=%s average_user_cpu=%s system_cpu=%s average_system_cpu=%s memory=%.1fMB average_memory=%.1fMB lock_held(read/write)=%s/%s tables=%d", report.Commands, report.Lapse, report.AverageLapse, report.UserCPU, report.AverageUserCPU, report.SystemCPU, report.AverageSystemCPU, report.MemoryMB, report.AverageMemoryMB, report.ReadLockHeld, report.WriteLockHeld, len(report.Tables)) } func (t trackingTotals) report() trackingReport { report := trackingReport{ Commands: t.commands, Lapse: t.lapse.String(), UserCPU: t.userCPU.String(), SystemCPU: t.systemCPU.String(), MemoryMB: t.memoryMB, ReadLockHeld: t.readLockHeld.String(), WriteLockHeld: t.writeLockHeld.String(), Tables: make([]tableLockReport, 0, len(t.tableLockHeld)), } if t.commands > 0 { commands := time.Duration(t.commands) report.AverageLapse = formatAverageMilliseconds(t.lapse / commands) report.AverageUserCPU = formatAverageMilliseconds(t.userCPU / commands) report.AverageSystemCPU = formatAverageMilliseconds(t.systemCPU / commands) report.AverageMemoryMB = t.memoryMB / float64(t.commands) } for table, held := range t.tableLockHeld { report.Tables = append(report.Tables, tableLockReport{ Table: table, ReadHeld: held.read.String(), WriteHeld: held.write.String(), }) } sort.Slice(report.Tables, func(i, j int) bool { return report.Tables[i].Table < report.Tables[j].Table }) return report } func (t trackingTotals) log(logger *logrus.Logger) { report := t.report() logger.Info(t.String()) for _, table := range report.Tables { logger.Infof("P4 table lock totals: table=%s read_held=%s write_held=%s", table.Table, table.ReadHeld, table.WriteHeld) } } func parseDate(value string) (time.Time, error) { return time.ParseInLocation("2006-01-02", value, time.Local) } func taggedRecords(output []byte, firstField string) []map[string]string { var records []map[string]string current := make(map[string]string) flush := func() { if len(current) > 0 { records = append(records, current) current = make(map[string]string) } } for _, line := range strings.Split(string(output), "\n") { line = strings.TrimSpace(line) if line == "" { flush() continue } if !strings.HasPrefix(line, "... ") { continue } parts := strings.SplitN(strings.TrimPrefix(line, "... "), " ", 2) if len(parts) != 2 { continue } field, value := parts[0], parts[1] if field == firstField && current[firstField] != "" { flush() } current[field] = value } flush() return records } func (p *p4Runner) localOrStreamDepots() ([]depotInfo, error) { output, err := p.run("depots") if err != nil { return nil, err } depots := make([]depotInfo, 0) for _, depot := range taggedRecords(output, "name") { if depot["type"] == "local" || depot["type"] == "stream" { info := depotInfo{name: depot["name"], typev: depot["type"]} if info.typev == "stream" { depth, err := parseStreamDepth(depot["depth"]) if err != nil { depth, err = p.streamDepthFromSpec(info.name) if err != nil { return nil, err } } info.streamDepth = depth } depots = append(depots, info) } } p.logger.Debugf("Found %d local or stream depots", len(depots)) return depots, nil } func (p *p4Runner) streamDepthFromSpec(depotName string) (int, error) { output, err := p.run("depot", "-o", depotName) if err != nil { return 0, err } for _, line := range strings.Split(string(output), "\n") { line = strings.TrimSpace(line) if strings.HasPrefix(line, "... StreamDepth ") { depth, err := parseStreamDepth(strings.TrimPrefix(line, "... StreamDepth ")) if err == nil { return depth, nil } } } return 0, fmt.Errorf("invalid stream depth for depot %q", depotName) } func parseStreamDepth(value string) (int, error) { depthText := value if index := strings.LastIndex(value, "/"); index >= 0 { depthText = value[index+1:] } depth, err := strconv.Atoi(depthText) if err != nil || depth < 1 { return 0, fmt.Errorf("invalid stream depth %q", value) } return depth, nil } func (p *p4Runner) changelistsByMonth(start, end time.Time) (map[string][]int, error) { output, err := p.run("changes", "-s", "submitted", "-m", "0") if err != nil { return nil, err } grouped := make(map[string][]int) for _, change := range taggedRecords(output, "change") { changeID := change["change"] timestamp := change["time"] if !completeChangeRecord(change) { p.logger.Debugf("Skipping incomplete tagged changelist record: change=%q time=%q", changeID, timestamp) continue } unixTime, err := strconv.ParseInt(timestamp, 10, 64) if err != nil { return nil, fmt.Errorf("parse change %q timestamp %q: %w", changeID, timestamp, err) } changeDate := time.Unix(unixTime, 0).In(time.Local) if changeDate.Before(start) || changeDate.After(end.AddDate(0, 0, 1).Add(-time.Nanosecond)) { continue } changeNumber, err := strconv.Atoi(changeID) if err != nil { return nil, fmt.Errorf("parse changelist number %q: %w", changeID, err) } month := changeDate.Format("2006-01") grouped[month] = append(grouped[month], changeNumber) } return grouped, nil } func completeChangeRecord(change map[string]string) bool { return change["change"] != "" && change["time"] != "" } func (p *p4Runner) filespecsForDepot(depot depotInfo, firstChange, lastChange int) ([]string, error) { revisionRange := fmt.Sprintf("@%d,%d", firstChange, lastChange) if depot.typev != "stream" { return []string{fmt.Sprintf("//%s/...%s", depot.name, revisionRange)}, nil } streams, found := p.streamDirectories[depot.name] if !found { streamDirsSpec, err := streamDirsSpec(depot) if err != nil { return nil, err } output, err := p.run("dirs", "-D", streamDirsSpec) if err != nil { return nil, err } streams = streamDirectories(output) if p.streamDirectories == nil { p.streamDirectories = make(map[string][]string) } p.streamDirectories[depot.name] = streams p.logger.Debugf("Cached %d streams in depot %s", len(streams), depot.name) } filespecs := make([]string, 0, len(streams)) for _, stream := range streams { filespecs = append(filespecs, stream+"/..."+revisionRange) } return filespecs, nil } func streamDirectories(output []byte) []string { streams := make([]string, 0) for _, directory := range taggedRecords(output, "dir") { stream := strings.TrimSuffix(directory["dir"], "/") if stream != "" { streams = append(streams, stream) } } return streams } func streamDirsSpec(depot depotInfo) (string, error) { if depot.typev != "stream" || depot.streamDepth < 1 { return "", fmt.Errorf("invalid stream depot %q depth %d", depot.name, depot.streamDepth) } return "//" + depot.name + strings.Repeat("/*", depot.streamDepth), nil } func (p *p4Runner) fileRevisionCounts(depots []depotInfo, firstChange, lastChange int) (revisionCounts, error) { counts := revisionCounts{} for _, depot := range depots { filespecs, err := p.filespecsForDepot(depot, firstChange, lastChange) if err != nil { return revisionCounts{}, err } for _, filespec := range filespecs { depotCounts, err := p.streamFileRevisionCounts(filespec) if err != nil { return revisionCounts{}, err } counts.addEdit += depotCounts.addEdit counts.deleted += depotCounts.deleted counts.branched += depotCounts.branched counts.purged += depotCounts.purged } } return counts, nil } func monthsBetween(start, end time.Time) []time.Time { month := time.Date(start.Year(), start.Month(), 1, 0, 0, 0, 0, time.Local) lastMonth := time.Date(end.Year(), end.Month(), 1, 0, 0, 0, 0, time.Local) months := make([]time.Time, 0) for !month.After(lastMonth) { months = append(months, month) month = month.AddDate(0, 1, 0) } return months } func (p *p4Runner) buildReport(start, end time.Time) ([]reportRow, error) { depots, err := p.localOrStreamDepots() if err != nil { return nil, err } changesByMonth, err := p.changelistsByMonth(start, end) if err != nil { return nil, err } rows := make([]reportRow, 0) for _, month := range monthsBetween(start, end) { key := month.Format("2006-01") changes := changesByMonth[key] row := reportRow{Month: key, Changes: len(changes)} if len(changes) > 0 { sort.Ints(changes) counts, err := p.fileRevisionCounts(depots, changes[0], changes[len(changes)-1]) if err != nil { return nil, fmt.Errorf("%s: %w", key, err) } row.AddEdit = counts.addEdit row.Deleted = counts.deleted row.Branched = counts.branched row.Purged = counts.purged } p.logger.Debugf("%s: %d submitted changes", key, row.Changes) rows = append(rows, row) } return rows, nil } func writeReport(rows []reportRow, output io.Writer) error { writer := csv.NewWriter(output) if err := writer.Write([]string{"month", "changes", "add_edit", "deleted", "branched", "purged"}); err != nil { return err } for _, row := range rows { if err := writer.Write([]string{ row.Month, strconv.Itoa(row.Changes), strconv.Itoa(row.AddEdit), strconv.Itoa(row.Deleted), strconv.Itoa(row.Branched), strconv.Itoa(row.Purged), }); err != nil { return err } } writer.Flush() return writer.Error() } func vmagentSettings() (endpoint, customer, password string, enabled bool, err error) { envData, envErr := os.ReadFile(vmagentEnvFile) passwordData, passwordErr := os.ReadFile(vmagentPasswordFile) if errors.Is(envErr, os.ErrNotExist) || errors.Is(passwordErr, os.ErrNotExist) { return "", "", "", false, nil } if envErr != nil { return "", "", "", false, envErr } if passwordErr != nil { return "", "", "", false, passwordErr } values := make(map[string]string) for _, line := range strings.Split(string(envData), "\n") { line = strings.TrimSpace(line) if line == "" || strings.HasPrefix(line, "#") { continue } key, value, found := strings.Cut(line, "=") if found { values[strings.TrimSpace(key)] = strings.TrimSpace(value) } } host := strings.TrimRight(values["VM_METRICS_HOST"], "/") customer = values["VM_CUSTOMER"] password = strings.TrimSpace(string(passwordData)) if host == "" || customer == "" || password == "" { return "", "", "", false, nil } return strings.Replace(host, ":9093", ":9092", 1) + "/json/", customer, password, true, nil } func pushJSONReport(rows []reportRow, tracking trackingTotals, instance string, timeout time.Duration, logger *logrus.Logger) error { endpoint, customer, password, enabled, err := vmagentSettings() if err != nil || !enabled { return err } query := url.Values{"customer": {customer}, "instance": {instance}} payload, err := dataPushPayload(rows, tracking) if err != nil { return err } logger.Debugf("Posting report to %s", endpoint) request, err := http.NewRequest(http.MethodPost, endpoint+"?"+query.Encode(), bytes.NewReader(payload)) if err != nil { return err } request.Header.Set("Authorization", "Basic "+base64.StdEncoding.EncodeToString([]byte(customer+":"+password))) request.Header.Set("Content-Type", "application/json") response, err := (&http.Client{Timeout: timeout}).Do(request) if err != nil { return err } defer response.Body.Close() if response.StatusCode < http.StatusOK || response.StatusCode >= http.StatusMultipleChoices { body, readErr := io.ReadAll(io.LimitReader(response.Body, 4096)) if readErr != nil { return fmt.Errorf("vmagent report push failed with HTTP %d", response.StatusCode) } if message := strings.TrimSpace(string(body)); message != "" { return fmt.Errorf("vmagent report push failed with HTTP %d: %s", response.StatusCode, message) } return fmt.Errorf("vmagent report push failed with HTTP %d", response.StatusCode) } return nil } func dataPushPayload(rows []reportRow, tracking trackingTotals) ([]byte, error) { report, err := json.Marshal(remoteReport{Version: Version, Rows: rows, Tracking: tracking.report()}) if err != nil { return nil, err } return json.Marshal([]map[string]string{{ "command": "p4changestats", "description": "Monthly Perforce submitted changelist and file revision statistics", "output": base64.StdEncoding.EncodeToString(report), "monitor_tag": "changestats", }}) } func mustJSON(rows []reportRow) []byte { data, err := json.Marshal(rows) if err != nil { panic(err) } return data } func run(args []string, stdout, stderr io.Writer) error { app := kingpin.New("p4changestats", description) app.ErrorWriter(stderr) app.UsageWriter(stderr) port := app.Flag("port", "P4PORT override").String() user := app.Flag("user", "P4USER override").String() client := app.Flag("client", "P4CLIENT override").String() p4bin := app.Flag("p4bin", "Path to p4 executable").Default("p4").String() outputPath := app.Flag("output", "CSV output file (default stdout)").String() instance := app.Flag("instance", "Data PushGateway instance name").Default(defaultInstance()).String() pushTimeoutValue := app.Flag("gateway-timeout", "Data PushGateway request timeout").Default(defaultPushTimeout.String()).String() debug := app.Flag("debug", "Print P4 commands and progress to stderr").Bool() startValue := app.Arg("start", "First included date (YYYY-MM-DD)").Required().String() endValue := app.Arg("end", "Last included date (YYYY-MM-DD)").Required().String() app.HelpFlag.Short('h') app.Version(Version) app.VersionFlag.Short('V') if _, err := app.Parse(args); err != nil { return err } start, err := parseDate(*startValue) if err != nil { return fmt.Errorf("invalid start date: %w", err) } end, err := parseDate(*endValue) if err != nil { return fmt.Errorf("invalid end date: %w", err) } if start.After(end) { return errors.New("start date must not be later than end date") } pushTimeout, err := time.ParseDuration(*pushTimeoutValue) if err != nil || pushTimeout <= 0 { return fmt.Errorf("invalid push timeout %q", *pushTimeoutValue) } logger := logrus.New() logger.SetOutput(stderr) logger.SetLevel(logrus.InfoLevel) if *debug { logger.SetLevel(logrus.DebugLevel) } logger.WithFields(logrus.Fields{ "start": *startValue, "end": *endValue, "port": *port, "user": *user, }).Debug("Building changelist statistics report") p4 := p4Runner{bin: *p4bin, port: *port, user: *user, client: *client, logger: logger} rows, err := p4.buildReport(start, end) if err != nil { return err } if *outputPath == "" { if err := writeReport(rows, stdout); err != nil { return err } } else { if err := os.MkdirAll(filepath.Dir(*outputPath), 0755); err != nil && filepath.Dir(*outputPath) != "." { return err } output, err := os.Create(*outputPath) if err != nil { return err } writeErr := writeReport(rows, output) closeErr := output.Close() if writeErr != nil { return writeErr } if closeErr != nil { return closeErr } } p4.tracking.log(logger) if err := pushJSONReport(rows, p4.tracking, *instance, pushTimeout, logger); err != nil { fmt.Fprintf(stderr, "%v\n", err) } return nil } func defaultInstance() string { hostname, err := os.Hostname() if err != nil { return "" } return strings.SplitN(hostname, ".", 2)[0] } func main() { if err := run(os.Args[1:], os.Stdout, os.Stderr); err != nil { fmt.Fprintln(os.Stderr, err) os.Exit(1) } }