Skip to content

Commit ba2a1bf

Browse files
committed
cli: add shard-index summary command
* add `ais bucket shard-index summary` with prefix, refresh, dont-wait, UUID polling support * render TAR object/shard coverage totals from the bucket shard-summary API * add CLI e2e coverage, JOB_ID validation, and CLI docs Signed-off-by: Tony Chen <a122774007@gmail.com>
1 parent b7696ee commit ba2a1bf

13 files changed

Lines changed: 304 additions & 33 deletions

File tree

cmd/cli/cli/bucket_hdlr.go

Lines changed: 34 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -153,12 +153,12 @@ const rechunkUsage = "Re-chunk bucket objects based on size threshold.\n" +
153153
indent1 + "\tBy default, rechunk operates only on in-cluster (cached) objects; use --sync-remote to also update remote backend."
154154

155155
// ais bucket shard-index
156-
// Parent command groups the shard-index lifecycle (build today; rm/show to follow in separate commits).
156+
// Parent command groups the shard-index lifecycle.
157+
// TODO: add shard-index rm subcommand to remove existing shard indexes.
157158
const shardIndexUsage = "Manage TAR shard indexes for fast random access into archives.\n" +
158159
indent1 + "\tSubcommands:\n" +
159-
indent1 + "\t- build\t- build a shard index for each TAR object in a bucket (see 'ais bucket shard-index build --help');\n" +
160-
indent1 + "\t- rm\t- (TODO, separate commit) remove existing shard indexes;\n" +
161-
indent1 + "\t- show\t- (TODO, separate commit) list/inspect existing shard indexes."
160+
indent1 + "\t- build\t- build a shard index for each TAR object in a bucket;\n" +
161+
indent1 + "\t- summary\t- summarize TAR objects and their shard-index coverage."
162162

163163
// ais bucket shard-index build
164164
const shardIndexBuildUsage = "Build a shard index for each TAR object in a bucket (for fast random access into archives).\n" +
@@ -170,6 +170,15 @@ const shardIndexBuildUsage = "Build a shard index for each TAR object in a bucke
170170
indent1 + "\t- 'ais bucket shard-index build ais://nnn --skip-verify'\t- fast re-run: trust existing indexes without re-verifying;\n" +
171171
indent1 + "\t- 'ais bucket shard-index build ais://nnn --wait'\t- start and wait for the job to finish."
172172

173+
// ais bucket shard-index summary
174+
const shardIndexSummaryUsage = "Summarize TAR objects in a bucket and their shard-index coverage.\n" +
175+
indent1 + "\tOnly local in-cluster objects are summarized.\n" +
176+
indent1 + "e.g.:\n" +
177+
indent1 + "\t- 'ais bucket shard-index summary ais://nnn'\t- summarize all TAR shards in 'ais://nnn';\n" +
178+
indent1 + "\t- 'ais bucket shard-index summary ais://nnn --prefix shards/'\t- summarize only TAR shards under 'shards/';\n" +
179+
indent1 + "\t- 'ais bucket shard-index summary ais://nnn/shards/'\t- same as above;\n" +
180+
indent1 + "\t- 'ais bucket shard-index summary ais://nnn --refresh 1s'\t- print periodic progress while summarizing."
181+
173182
// ais bucket ... props
174183
const setBpropsUsage = "Update bucket properties; the command accepts both JSON-formatted input and plain Name=Value pairs,\n" +
175184
indent1 + "\te.g.:\n" +
@@ -318,6 +327,13 @@ var (
318327
waitJobXactFinishedFlag,
319328
nonverboseFlag,
320329
},
330+
cmdSummary: {
331+
verbObjPrefixFlag,
332+
refreshFlag,
333+
unitsFlag,
334+
dontWaitFlag,
335+
noHeaderFlag,
336+
},
321337
}
322338
)
323339

@@ -403,8 +419,14 @@ var (
403419
Flags: sortFlags(bucketCmdsFlags[cmdShardIndexBuild]),
404420
Action: shardIndexBuildHandler,
405421
},
406-
// TODO (separate commit): {Name: "rm", Action: shardIndexRmHandler} - remove existing shard indexes
407-
// TODO (separate commit): {Name: "show", Action: shardIndexShowHandler} - list/inspect existing shard indexes
422+
{
423+
Name: cmdSummary,
424+
Usage: shardIndexSummaryUsage,
425+
ArgsUsage: bucketArgument + " [JOB_ID]",
426+
Flags: sortFlags(bucketCmdsFlags[cmdSummary]),
427+
Action: shardIndexSummaryHandler,
428+
BashComplete: bucketCompletions(bcmplop{}),
429+
},
408430
},
409431
},
410432
makeAlias(&showCmdBucket, &mkaliasOpts{newName: commandShow}),
@@ -622,13 +644,9 @@ func rechunkBucketHandler(c *cli.Context) error {
622644
return err
623645
}
624646

625-
prefix := parseStrFlag(c, verbObjPrefixFlag)
626-
if objName != "" && prefix != "" && !strings.HasPrefix(prefix, objName) {
627-
return fmt.Errorf("cannot handle embedded prefix ('%s') and --prefix flag ('%s') simultaneously - prefix flag must start with embedded prefix",
628-
objName, prefix)
629-
}
630-
if prefix == "" {
631-
prefix = objName
647+
prefix, err := parseBckObjPrefix(c, objName)
648+
if err != nil {
649+
return err
632650
}
633651

634652
// Start rechunk
@@ -728,14 +746,9 @@ func shardIndexBuildHandler(c *cli.Context) error {
728746
return err
729747
}
730748

731-
// combine embedded prefix (if any) with --prefix flag
732-
prefix := parseStrFlag(c, verbObjPrefixFlag)
733-
if objName != "" && prefix != "" && !strings.HasPrefix(prefix, objName) {
734-
return fmt.Errorf("conflict between embedded prefix (%q) and --prefix flag (%q): --prefix must start with the embedded prefix",
735-
objName, prefix)
736-
}
737-
if prefix == "" {
738-
prefix = objName
749+
prefix, err := parseBckObjPrefix(c, objName)
750+
if err != nil {
751+
return err
739752
}
740753

741754
msg := &apc.IndexShardMsg{

cmd/cli/cli/const.go

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -117,11 +117,10 @@ const (
117117
commandRechunk = apc.ActRechunk
118118
commandShardIndex = "shard-index" // parent; xaction kind is apc.ActIndexShard
119119
cmdShardIndexBuild = "build"
120-
// TODO: cmdShardIndexRm = "rm" - remove existing shard indexes
121-
// TODO: cmdShardIndexShow = "show" - list/inspect existing shard indexes
120+
// TODO: cmdShardIndexRm = "rm" - remove existing shard indexes
122121
cmdStgCleanup = "cleanup" // display name for apc.ActStoreCleanup
123122
cmdScrub = "validate"
124-
cmdSummary = "summary" // ditto apc.ActSummaryBck
123+
cmdSummary = "summary"
125124

126125
cmdCluster = commandCluster
127126
cmdDashboard = commandDashboard

cmd/cli/cli/parse_uri.go

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -189,6 +189,18 @@ func parseBckObjURI(c *cli.Context, uri string, emptyObjnameOK bool) (bck cmn.Bc
189189
return bck, objName, nil
190190
}
191191

192+
func parseBckObjPrefix(c *cli.Context, objName string) (string, error) {
193+
prefix := parseStrFlag(c, verbObjPrefixFlag)
194+
if objName != "" && prefix != "" && !strings.HasPrefix(prefix, objName) {
195+
return "", fmt.Errorf("cannot handle embedded prefix ('%s') and --prefix flag ('%s') simultaneously - prefix flag must start with embedded prefix",
196+
objName, prefix)
197+
}
198+
if prefix == "" {
199+
prefix = objName
200+
}
201+
return prefix, nil
202+
}
203+
192204
//
193205
// - handle (obj names) list, template (range), embedded prefix, and single object name
194206
// - possibly call list-objects (via `lsObjVsPref`) to disambiguate

cmd/cli/cli/shard_idx_hdlr.go

Lines changed: 179 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,179 @@
1+
// Package cli provides easy-to-use commands to manage, monitor, and utilize AIS clusters.
2+
/*
3+
* Copyright (c) 2026, NVIDIA CORPORATION. All rights reserved.
4+
*/
5+
package cli
6+
7+
import (
8+
"fmt"
9+
"net/http"
10+
"strings"
11+
12+
"github.com/NVIDIA/aistore/api"
13+
"github.com/NVIDIA/aistore/api/apc"
14+
"github.com/NVIDIA/aistore/cmd/cli/teb"
15+
"github.com/NVIDIA/aistore/cmn"
16+
"github.com/NVIDIA/aistore/cmn/cos"
17+
18+
"github.com/urfave/cli"
19+
)
20+
21+
type (
22+
shardSummCtx struct {
23+
c *cli.Context
24+
units string
25+
bck cmn.Bck
26+
prefix string
27+
msg apc.ShardSummMsg
28+
args api.ShardSummArgs
29+
l int
30+
n int
31+
}
32+
shardSummRow struct {
33+
Name string
34+
NotIndexed uint64
35+
apc.ShardSummResult
36+
}
37+
)
38+
39+
func shardIndexSummaryHandler(c *cli.Context) error {
40+
if c.NArg() == 0 {
41+
return incorrectUsageMsg(c, "missing bucket name")
42+
}
43+
if c.NArg() > 2 {
44+
return incorrectUsageMsg(c, "", c.Args()[2:])
45+
}
46+
47+
bck, objName, err := parseBckObjURI(c, c.Args().Get(0), true /*optObjName*/)
48+
if err != nil {
49+
return err
50+
}
51+
prefix, err := parseBckObjPrefix(c, objName)
52+
if err != nil {
53+
return err
54+
}
55+
ctx, err := newShardSummCtxMsg(c, bck, prefix)
56+
if err != nil {
57+
return err
58+
}
59+
60+
// Optional JOB_ID means polling an existing summary job.
61+
isNew := true
62+
if xid := c.Args().Get(1); xid != "" {
63+
if !cos.IsValidUUID(xid) {
64+
return incorrectUsageMsg(c, "invalid job ID %q", xid)
65+
}
66+
if !ctx.args.DontWait {
67+
return incorrectUsageMsg(c, "JOB_ID requires %s", flprn(dontWaitFlag))
68+
}
69+
ctx.msg.UUID = xid
70+
isNew = false
71+
}
72+
73+
// Start or poll the summary, depending on msg.UUID and --dont-wait.
74+
xid, res, err := api.GetBucketShardSummary(apiBP, ctx.bck, &ctx.msg, ctx.args)
75+
76+
// New --dont-wait starts the job; if no snapshot exists yet, print the xaction ID.
77+
dontWait := flagIsSet(c, dontWaitFlag)
78+
if err == nil && dontWait && isNew && res.IsEmpty() {
79+
actionDone(c, shardSummStartedMsg(xid, bck.Cname(prefix), isNew))
80+
return nil
81+
}
82+
83+
// For --dont-wait polling, 202 means no snapshot yet; 206 prints the latest partial snapshot.
84+
var status int
85+
if err != nil {
86+
if herr := cmn.AsErrHTTP(err); herr != nil {
87+
status = herr.Status
88+
}
89+
if dontWait && status == http.StatusAccepted {
90+
actionDone(c, shardSummStartedMsg(xid, bck.Cname(prefix), isNew))
91+
return nil
92+
}
93+
if dontWait && status == http.StatusPartialContent {
94+
msg := fmt.Sprintf("%s[%s] is still running - showing partial results:", apc.ActSummaryShard, ctx.msg.UUID)
95+
actionNote(c, msg)
96+
err = nil
97+
}
98+
}
99+
if err != nil {
100+
return V(err)
101+
}
102+
// Print the summary table.
103+
return ctx.print(res)
104+
}
105+
106+
func shardSummStartedMsg(xid, bucket string, isNew bool) string {
107+
verb := "has started"
108+
if !isNew {
109+
verb = "is running"
110+
}
111+
return fmt.Sprintf("Job %s[%s] %s. To monitor, run 'ais bucket shard-index summary %s %s %s'",
112+
apc.ActSummaryShard, xid, verb, bucket, xid, flprn(dontWaitFlag))
113+
}
114+
115+
func newShardSummCtxMsg(c *cli.Context, bck cmn.Bck, prefix string) (*shardSummCtx, error) {
116+
units, err := parseUnitsFlag(c, unitsFlag)
117+
if err != nil {
118+
return nil, err
119+
}
120+
ctx := &shardSummCtx{
121+
c: c,
122+
units: units,
123+
bck: bck,
124+
prefix: prefix,
125+
}
126+
ctx.msg.Prefix = prefix
127+
if ctx.args.DontWait = flagIsSet(c, dontWaitFlag); ctx.args.DontWait {
128+
return ctx, nil
129+
}
130+
131+
ctx.args.CallAfter = _refreshRate(c)
132+
ctx.args.Callback = ctx.progress
133+
return ctx, nil
134+
}
135+
136+
func (ctx *shardSummCtx) progress(res *apc.ShardSummResult, done bool) {
137+
if done {
138+
if ctx.n > 0 {
139+
fmt.Fprintln(ctx.c.App.Writer)
140+
}
141+
return
142+
}
143+
if res == nil || res.IsEmpty() {
144+
return
145+
}
146+
ctx.n++
147+
148+
s := fmt.Sprintf("%s: %s/%s indexed (%s indexed, %s total)",
149+
ctx.bck.Cname(ctx.prefix),
150+
cos.FormatBigI64(int64(res.Shards)),
151+
cos.FormatBigI64(int64(res.TarObjs)),
152+
teb.FmtSize(int64(res.ShardSize), ctx.units, 2),
153+
teb.FmtSize(int64(res.TarSize), ctx.units, 2))
154+
if ctx.l < len(s) {
155+
ctx.l = len(s) + 4
156+
}
157+
s += strings.Repeat(" ", ctx.l-len(s))
158+
fmt.Fprintf(ctx.c.App.Writer, "\r%s", s)
159+
}
160+
161+
func (ctx *shardSummCtx) print(res *apc.ShardSummResult) error {
162+
row := shardSummRow{
163+
Name: ctx.bck.Cname(ctx.prefix),
164+
NotIndexed: shardSummNotIndexed(res),
165+
ShardSummResult: *res,
166+
}
167+
opts := teb.Opts{AltMap: teb.FuncMapUnits(ctx.units, false /*incl. calendar date*/)}
168+
if flagIsSet(ctx.c, noHeaderFlag) {
169+
return teb.Print([]shardSummRow{row}, teb.ShardSummariesBody, opts)
170+
}
171+
return teb.Print([]shardSummRow{row}, teb.ShardSummariesTmpl, opts)
172+
}
173+
174+
func shardSummNotIndexed(res *apc.ShardSummResult) uint64 {
175+
if res.Shards > res.TarObjs {
176+
return 0
177+
}
178+
return res.TarObjs - res.Shards
179+
}

cmd/cli/go.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ module github.com/NVIDIA/aistore/cmd/cli
33
go 1.26
44

55
require (
6-
github.com/NVIDIA/aistore v1.4.7-0.20260603190201-06f8bfc83ff0
6+
github.com/NVIDIA/aistore v1.4.8-0.20260610215907-6394b108d4fc
77
github.com/fatih/color v1.19.0
88
github.com/json-iterator/go v1.1.12
99
github.com/onsi/ginkgo/v2 v2.28.1

cmd/cli/go.sum

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@ github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1
33
github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM=
44
github.com/NVIDIA/aistore v1.4.7-0.20260603190201-06f8bfc83ff0 h1:7AhTtDb1Enmn2m0zWK4NkTWZa+cXlYBPsSa41T64CtU=
55
github.com/NVIDIA/aistore v1.4.7-0.20260603190201-06f8bfc83ff0/go.mod h1:gQVg+bfLzebJFLhHTQ3cqyLI1EdVm32LzF1f1ETDYgc=
6+
github.com/NVIDIA/aistore v1.4.8-0.20260610215907-6394b108d4fc h1:RpA66P6pBoJnpXXZs0mF3Z2vSUHhS1QX31ysrDKHMMI=
7+
github.com/NVIDIA/aistore v1.4.8-0.20260610215907-6394b108d4fc/go.mod h1:gQVg+bfLzebJFLhHTQ3cqyLI1EdVm32LzF1f1ETDYgc=
68
github.com/OneOfOne/xxhash v1.2.8 h1:31czK/TI9sNkxIKfaUfGlU47BAxQ0ztGgd9vPyqimf8=
79
github.com/OneOfOne/xxhash v1.2.8/go.mod h1:eZbhyaAYD41SGSSsnmcpxVoRiQ/MPUTjUdIIOT9Um7Q=
810
github.com/VividCortex/ewma v1.1.1/go.mod h1:2Tkkvm3sRDVXaiyucHiACn4cqf7DpdyLvmxzcbUokwA=

cmd/cli/teb/templates.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -350,6 +350,15 @@ const (
350350
"{{FormatBytesUns $v.TotalSize.PresentObjs 2}} {{FormatBytesUns $v.TotalSize.RemoteObjs 2}}\t {{$v.UsedPct}}%\n" +
351351
"{{end}}"
352352

353+
// Shard index summary templates
354+
ShardSummariesTmpl = "BUCKET\t TAR OBJECTS\t TAR SIZE\t SHARDS\t SHARD SIZE\t NOT INDEXED\t ARCHIVED OBJECTS\t STALE\t INVALID\n" +
355+
ShardSummariesBody
356+
ShardSummariesBody = "{{range $v := . }}" +
357+
"{{$v.Name}}\t {{$v.TarObjs}}\t {{FormatBytesUns $v.TarSize 2}}\t {{$v.Shards}}\t " +
358+
"{{FormatBytesUns $v.ShardSize 2}}\t {{$v.NotIndexed}}\t {{$v.ArchivedObjs}}\t " +
359+
"{{$v.StaleIndexes}}\t {{$v.InvalidIndexes}}\n" +
360+
"{{end}}"
361+
353362
// For `object put` mass uploader. A caller adds to the template
354363
// total count and size. That is why the template ends with \t
355364
MultiPutTmpl = "Files to upload:\nEXTENSION\t COUNT\t SIZE\n" +

cmd/cli/test/index_shard.in

Lines changed: 22 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,21 +1,38 @@
1-
# Test shard-index build command
1+
# Test shard-index build and summary commands
22

33
# Create bucket
44
ais bucket create ais://$BUCKET // IGNORE
55

6-
# Generate 5 TAR shards (2 files each) directly into the bucket
7-
ais archive gen-shards "ais://$BUCKET/shard-{001..005}.tar" --fcount 2 // IGNORE
6+
# Generate 5 TAR shards (2 files each) under the tested prefix
7+
ais archive gen-shards "ais://$BUCKET/shards/shard-{001..005}.tar" --fcount 2 // IGNORE
8+
9+
# Generate one TAR shard outside the tested prefix
10+
ais archive gen-shards "ais://$BUCKET/skip/shard-{001..001}.tar" --fcount 2 // IGNORE
811

912
# Add a non-TAR object (must be skipped, not indexed)
1013
head -c 1024 /dev/urandom | ais object put - ais://$BUCKET/readme.md // IGNORE
1114

15+
# Reject invalid job IDs instead of starting a new summary job
16+
ais bucket shard-index summary ais://$BUCKET bad-id --dont-wait // FAIL "invalid job ID"
17+
ais bucket shard-index summary ais://$BUCKET abcDEF123 // FAIL "JOB_ID requires --dont-wait"
18+
19+
# Before indexing, the prefix summary must find TAR objects but no indexed shards
20+
ais bucket shard-index summary ais://$BUCKET/shards/ --no-headers | awk '{print "pre-summary=" $2 "," $4 "," $6 "," $7}'
21+
1222
# Run shard-index build and wait for completion
13-
ais bucket shard-index build ais://$BUCKET --wait // IGNORE
23+
ais bucket shard-index build ais://$BUCKET/shards/ --wait // IGNORE
1424

1525
# Extract processed-object count from the job summary;
16-
# 5 TAR shards were indexed, readme.md was skipped (non-TAR) --> expect indexed=5
26+
# 5 TAR shards under shards/ were indexed, skip/ and readme.md were not selected --> expect indexed=5
1727
# Sum the OBJECTS column across all target rows.
1828
ais show job index-shard ais://$BUCKET --all --no-headers 2>&1 | awk '$3 == "index-shard" {sum += $5} END {print "indexed=" sum}'
1929

30+
# After indexing, prefix summary must count only shards/ TAR objects
31+
ais bucket shard-index summary ais://$BUCKET --prefix shards/ --no-headers | awk '{print "summary=" $2 "," $4 "," $6 "," $7}'
32+
33+
# Start asynchronously and poll by UUID
34+
ais bucket shard-index summary ais://$BUCKET/shards/ --dont-wait | awk -F'[][]' '{print $2}' // SAVE_RESULT
35+
for i in $(seq 1 30); do out=$(ais bucket shard-index summary ais://$BUCKET/shards/ $RESULT --dont-wait --no-headers 2>/dev/null || true); if echo "$out" | grep -q "^ais://"; then echo "$out" | awk '{print "poll-summary=" $2 "," $4 "," $6 "," $7}'; break; fi; sleep 0.1; done
36+
2037
# Clean up
2138
ais bucket rm --yes ais://$BUCKET // IGNORE

cmd/cli/test/index_shard.stdout

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1 +1,4 @@
1+
pre-summary=5,0,5,0
12
indexed=5
3+
summary=5,5,0,10
4+
poll-summary=5,5,0,10

0 commit comments

Comments
 (0)