// 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=<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)
}
}