Skip to content

Commit 21ff893

Browse files
Support async DB listing & job retrieval
Add higher-level client helpers and refactor callers: implement Client.GetJob to handle both wrapped and flat job responses, and Client.ListDatabases to handle synchronous (200) and asynchronous (202 -> job) /api/dbs responses. Introduce databasesFromJobResult to extract database rows from job results and update WaitForJob to use GetJob. Update cmd/db.go and cmd/jobs.go to call the new client methods and simplify response handling and output. These changes centralize API parsing and enable polling for async database listings.
1 parent 0a25af3 commit 21ff893

3 files changed

Lines changed: 128 additions & 14 deletions

File tree

cmd/db.go

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -24,35 +24,36 @@ var dbListCmd = &cobra.Command{
2424
return err
2525
}
2626

27-
var result struct {
28-
Data []map[string]interface{} `json:"data"`
29-
}
30-
31-
if err := client.Get("/api/dbs", &result); err != nil {
27+
dbs, err := client.ListDatabases()
28+
if err != nil {
3229
output.Error("Failed to list databases: %s", err)
3330
return err
3431
}
3532

33+
result := struct {
34+
Data []map[string]interface{} `json:"data"`
35+
}{Data: dbs}
36+
3637
if jsonFlag {
3738
output.PrintJSON(result)
3839
return nil
3940
}
4041

41-
if len(result.Data) == 0 {
42+
if len(dbs) == 0 {
4243
output.Warn("No databases found")
4344
return nil
4445
}
4546

4647
output.Header("Databases")
4748
t := output.NewTable("NAME", "SIZE")
48-
for _, db := range result.Data {
49+
for _, db := range dbs {
4950
t.Row(
5051
str(db, "name"),
5152
str(db, "size"),
5253
)
5354
}
5455
t.Flush()
55-
output.Dim.Printf(" Total: %d database(s)\n\n", len(result.Data))
56+
output.Dim.Printf(" Total: %d database(s)\n\n", len(dbs))
5657
return nil
5758
},
5859
}

cmd/jobs.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,8 +25,8 @@ var jobsShowCmd = &cobra.Command{
2525
return err
2626
}
2727

28-
var job api.JobStatus
29-
if err := client.Get(fmt.Sprintf("/api/jobs/%s", args[0]), &job); err != nil {
28+
job, err := client.GetJob(args[0])
29+
if err != nil {
3030
output.Error("Failed to get job: %s", err)
3131
return err
3232
}

internal/api/client.go

Lines changed: 117 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -140,6 +140,119 @@ func (c *Client) Get(path string, result interface{}) error {
140140
return c.parseResponse(resp, result)
141141
}
142142

143+
// GetJob fetches GET /api/jobs/{id}. The API returns either { "data": { job } } or a flat job object.
144+
func (c *Client) GetJob(jobID string) (*JobStatus, error) {
145+
resp, err := c.request("GET", fmt.Sprintf("/api/jobs/%s", jobID), nil)
146+
if err != nil {
147+
return nil, err
148+
}
149+
defer resp.Body.Close()
150+
data, err := io.ReadAll(resp.Body)
151+
if err != nil {
152+
return nil, err
153+
}
154+
if resp.StatusCode >= 400 {
155+
var apiErr APIError
156+
if json.Unmarshal(data, &apiErr) == nil && apiErr.Message != "" {
157+
return nil, &apiErr
158+
}
159+
return nil, fmt.Errorf("API error (HTTP %d): %s", resp.StatusCode, string(data))
160+
}
161+
var outer map[string]json.RawMessage
162+
if err := json.Unmarshal(data, &outer); err != nil {
163+
return nil, fmt.Errorf("parsing job response: %w", err)
164+
}
165+
if raw, ok := outer["data"]; ok {
166+
var job JobStatus
167+
if err := json.Unmarshal(raw, &job); err != nil {
168+
return nil, fmt.Errorf("parsing job data: %w", err)
169+
}
170+
return &job, nil
171+
}
172+
var job JobStatus
173+
if err := json.Unmarshal(data, &job); err != nil {
174+
return nil, fmt.Errorf("parsing job response: %w", err)
175+
}
176+
return &job, nil
177+
}
178+
179+
// ListDatabases returns GET /api/dbs. The server may answer synchronously (200 + data) or
180+
// asynchronously (202 + job_id); in the latter case the job is polled until completion.
181+
func (c *Client) ListDatabases() ([]map[string]interface{}, error) {
182+
resp, err := c.request("GET", "/api/dbs", nil)
183+
if err != nil {
184+
return nil, err
185+
}
186+
defer resp.Body.Close()
187+
body, err := io.ReadAll(resp.Body)
188+
if err != nil {
189+
return nil, err
190+
}
191+
if resp.StatusCode >= 400 {
192+
var apiErr APIError
193+
if json.Unmarshal(body, &apiErr) == nil && apiErr.Message != "" {
194+
return nil, &apiErr
195+
}
196+
return nil, fmt.Errorf("API error (HTTP %d): %s", resp.StatusCode, string(body))
197+
}
198+
199+
switch resp.StatusCode {
200+
case http.StatusOK:
201+
var v struct {
202+
Data []map[string]interface{} `json:"data"`
203+
}
204+
if err := json.Unmarshal(body, &v); err != nil {
205+
return nil, fmt.Errorf("parsing database list: %w", err)
206+
}
207+
if v.Data == nil {
208+
return []map[string]interface{}{}, nil
209+
}
210+
return v.Data, nil
211+
case http.StatusAccepted:
212+
var async AsyncResponse
213+
if err := json.Unmarshal(body, &async); err != nil {
214+
return nil, fmt.Errorf("parsing async response: %w", err)
215+
}
216+
if async.JobID == nil || async.JobIDString() == "" {
217+
return nil, fmt.Errorf("no job_id in database list response")
218+
}
219+
job, err := c.WaitForJob(async.JobIDString())
220+
if err != nil {
221+
return nil, err
222+
}
223+
return databasesFromJobResult(job.Result)
224+
default:
225+
return nil, fmt.Errorf("unexpected HTTP %d listing databases: %s", resp.StatusCode, string(body))
226+
}
227+
}
228+
229+
func databasesFromJobResult(result interface{}) ([]map[string]interface{}, error) {
230+
if result == nil {
231+
return []map[string]interface{}{}, nil
232+
}
233+
m, ok := result.(map[string]interface{})
234+
if !ok {
235+
return nil, fmt.Errorf("unexpected job result format for database list")
236+
}
237+
raw, ok := m["databases"]
238+
if !ok || raw == nil {
239+
return []map[string]interface{}{}, nil
240+
}
241+
arr, ok := raw.([]interface{})
242+
if !ok {
243+
return nil, fmt.Errorf("job result \"databases\" is not an array")
244+
}
245+
out := make([]map[string]interface{}, 0, len(arr))
246+
for _, item := range arr {
247+
row, ok := item.(map[string]interface{})
248+
if !ok {
249+
continue
250+
}
251+
out = append(out, row)
252+
}
253+
return out, nil
254+
}
255+
143256
func (c *Client) Post(path string, body interface{}, result interface{}) error {
144257
resp, err := c.request("POST", path, body)
145258
if err != nil {
@@ -191,23 +304,23 @@ func (c *Client) WaitForJob(jobID string) (*JobStatus, error) {
191304

192305
maxAttempts := 120
193306
for i := 0; i < maxAttempts; i++ {
194-
var job JobStatus
195-
if err := c.Get(fmt.Sprintf("/api/jobs/%s", jobID), &job); err != nil {
307+
job, err := c.GetJob(jobID)
308+
if err != nil {
196309
time.Sleep(2 * time.Second)
197310
continue
198311
}
199312

200313
switch job.Status {
201314
case "completed", "success", "finished":
202315
s.Stop()
203-
return &job, nil
316+
return job, nil
204317
case "failed", "error":
205318
s.Stop()
206319
errMsg := job.Error
207320
if errMsg == "" {
208321
errMsg = "job failed"
209322
}
210-
return &job, fmt.Errorf("%s", errMsg)
323+
return job, fmt.Errorf("%s", errMsg)
211324
}
212325

213326
interval := 2 * time.Second

0 commit comments

Comments
 (0)