Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,11 @@

<!-- Add manual release notes here. They will be merged into the generated changelog at release time. -->

### AI Task Builder

- Support export slicing: `aitaskbuilder batch export` and `collection export` accept `--study-id`, `--from`, and `--to` to narrow an export to a subset of responses
- Add `aitaskbuilder batch export list`/`delete`/`download` and `collection export list`/`delete`/`download` subcommands to list, remove, and (re-)download export jobs without starting a new export

## 1.2.5

- Maintenance and dependency updates
Expand Down
5 changes: 5 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -331,6 +331,8 @@ Operations are grouped as they appear in [`contract_test/contract_test.go`](cont
| `aiTaskBuilder_GetBatchSyncStatus` | GET | `/api/v1/data-collection/batches/{batch_id}/syncs/{sync_id}` | ✅ `GetAITaskBuilderBatchSyncStatus` |
| `aiTaskBuilder_RequestBatchExport` | POST | `/api/v1/data-collection/batches/{batch_id}/export` | ✅ `InitiateBatchExport` |
| `aiTaskBuilder_GetBatchExportStatus` | GET | `/api/v1/data-collection/batches/{batch_id}/export/{export_id}` | ✅ `GetBatchExportStatus` |
| `aiTaskBuilder_ListBatchExportJobs` | GET | `/api/v1/data-collection/batches/{batch_id}/export` | ✅ `ListBatchExportJobs` |
| `aiTaskBuilder_DeleteBatchExport` | DELETE | `/api/v1/data-collection/batches/{batch_id}/export/{export_id}` | ✅ `DeleteBatchExport` |

</details>

Expand Down Expand Up @@ -373,6 +375,8 @@ Operations are grouped as they appear in [`contract_test/contract_test.go`](cont
| `aiTaskBuilder_GetCollectionResponses` | GET | `/api/v1/data-collection/collections/{collection_id}/responses` | ➖ Not exposed in the CLI |
| `aiTaskBuilder_RequestCollectionExport` | POST | `/api/v1/data-collection/collections/{collection_id}/export` | ✅ `InitiateCollectionExport` |
| `aiTaskBuilder_GetCollectionExportStatus` | GET | `/api/v1/data-collection/collections/{collection_id}/export/{export_id}` | ✅ `GetCollectionExportStatus` |
| `aiTaskBuilder_ListCollectionExportJobs` | GET | `/api/v1/data-collection/collections/{collection_id}/export` | ✅ `ListCollectionExportJobs` |
| `aiTaskBuilder_DeleteCollectionExport` | DELETE | `/api/v1/data-collection/collections/{collection_id}/export/{export_id}` | ✅ `DeleteCollectionExport` |

</details>

Expand All @@ -396,6 +400,7 @@ Operations are grouped as they appear in [`contract_test/contract_test.go`](cont
| `messages_SendMessageToParticipantGroup` | POST | `/api/v1/messages/participant-group/` | ✅ `SendGroupMessage` |
| `messages_GetUnreadMessages` | GET | `/api/v1/messages/unread/` | ✅ `GetUnreadMessages` |
| `messages_GetConversations` | GET | `/api/v1/conversations/` | ➖ Not exposed in the CLI |
| `messages_CreateConversation` | POST | `/api/v1/conversations/` | ➖ Not exposed in the CLI |
| `messages_GetConversationMessages` | GET | `/api/v1/conversations/{conversation_id}/messages/` | ➖ Not exposed in the CLI |

</details>
Expand Down
91 changes: 85 additions & 6 deletions client/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,8 +60,10 @@ type API interface {

GetCollections(workspaceID string, limit, offset int) (*ListCollectionsResponse, error)
GetCollection(ID string) (*model.Collection, error)
InitiateCollectionExport(collectionID string) (*CollectionExportResponse, error)
InitiateCollectionExport(collectionID string, filter ExportFilter) (*CollectionExportResponse, error)
GetCollectionExportStatus(collectionID, exportID string) (*CollectionExportResponse, error)
ListCollectionExportJobs(collectionID string) ([]ExportJobListItem, error)
DeleteCollectionExport(collectionID, exportID string) error
UpdateCollection(ID string, collection model.UpdateCollection) (*model.Collection, error)

GetHooks(workspaceID string, enabled bool, limit, offset int) (*ListHooksResponse, error)
Expand Down Expand Up @@ -134,8 +136,10 @@ type API interface {
GetAITaskBuilderResponses(batchID string) (*GetAITaskBuilderResponsesResponse, error)
GetAITaskBuilderTasks(batchID string) (*GetAITaskBuilderTasksResponse, error)
GetAITaskBuilderTaskGroups(batchID string) (*GetAITaskBuilderTaskGroupsResponse, error)
InitiateBatchExport(batchID string) (*BatchExportResponse, error)
InitiateBatchExport(batchID string, filter ExportFilter) (*BatchExportResponse, error)
GetBatchExportStatus(batchID, exportID string) (*BatchExportResponse, error)
ListBatchExportJobs(batchID string) ([]ExportJobListItem, error)
DeleteBatchExport(batchID, exportID string) error
SyncAITaskBuilderBatch(batchID string) (*AITaskBuilderBatchSyncResponse, error)
GetAITaskBuilderBatchSyncStatus(batchID, syncID string) (*AITaskBuilderBatchSyncResponse, error)
GetAITaskBuilderDatasetStatus(datasetID string) (*GetAITaskBuilderDatasetStatusResponse, error)
Expand Down Expand Up @@ -571,10 +575,13 @@ func (c *Client) GetCollection(ID string) (*model.Collection, error) {
// InitiateCollectionExport starts a collection export job via POST.
// Returns "generating" + ExportID (202) if a new job was enqueued,
// or "complete" + URL immediately (200) if a valid export already exists.
func (c *Client) InitiateCollectionExport(collectionID string) (*CollectionExportResponse, error) {
// filter optionally slices the export down to responses for a specific
// study_id and/or a from/to created_at date range; a zero-value ExportFilter
// requests a full, unfiltered export.
func (c *Client) InitiateCollectionExport(collectionID string, filter ExportFilter) (*CollectionExportResponse, error) {
var response CollectionExportResponse

url := fmt.Sprintf("/api/v1/data-collection/collections/%s/export", collectionID)
url := fmt.Sprintf("/api/v1/data-collection/collections/%s/export%s", collectionID, filter.query())
_, err := c.ExecuteBuilder().PostRequest(url).Decode(&response).Execute()
if err != nil {
return nil, err
Expand All @@ -597,6 +604,38 @@ func (c *Client) GetCollectionExportStatus(collectionID, exportID string) (*Coll
return &response, nil
}

// ListCollectionExportJobs returns every export job requested for a
// collection, most recent first. Manually deleted jobs are excluded.
func (c *Client) ListCollectionExportJobs(collectionID string) ([]ExportJobListItem, error) {
var response []ExportJobListItem

url := fmt.Sprintf("/api/v1/data-collection/collections/%s/export", collectionID)
_, err := c.Execute(http.MethodGet, url, nil, &response)
if err != nil {
return nil, fmt.Errorf("unable to fulfil request %s: %s", url, err)
}

return response, nil
}

// DeleteCollectionExport permanently deletes an export job and, if it
// completed, its ZIP archive. An export that is still generating cannot be
// deleted (the API returns 409 Conflict).
func (c *Client) DeleteCollectionExport(collectionID, exportID string) error {
url := fmt.Sprintf("/api/v1/data-collection/collections/%s/export/%s", collectionID, exportID)
httpResponse, err := c.Execute(http.MethodDelete, url, nil, nil)
if err != nil {
return fmt.Errorf("unable to fulfil request %s: %s", url, err)
}

if httpResponse.StatusCode != http.StatusNoContent {
body, _ := io.ReadAll(httpResponse.Body)
return fmt.Errorf("unexpected status code %d: %s", httpResponse.StatusCode, string(body))
}

return nil
}

// UpdateStudy is responsible for updating the Study with a PATCH request.
func (c *Client) UpdateStudy(ID string, study any) (*model.Study, error) {
var response model.Study
Expand Down Expand Up @@ -1549,10 +1588,13 @@ func (c *Client) GetAITaskBuilderTaskGroups(batchID string) (*GetAITaskBuilderTa
// InitiateBatchExport starts a batch export job via POST.
// Returns "generating" + ExportID (202) if a new job was enqueued,
// or "complete" + URL immediately (200) if a valid export already exists.
func (c *Client) InitiateBatchExport(batchID string) (*BatchExportResponse, error) {
// filter optionally slices the export down to responses for a specific
// study_id and/or a from/to created_at date range; a zero-value ExportFilter
// requests a full, unfiltered export.
func (c *Client) InitiateBatchExport(batchID string, filter ExportFilter) (*BatchExportResponse, error) {
var response BatchExportResponse

url := fmt.Sprintf("/api/v1/data-collection/batches/%s/export", batchID)
url := fmt.Sprintf("/api/v1/data-collection/batches/%s/export%s", batchID, filter.query())
_, err := c.ExecuteBuilder().
PostRequest(url).
Status(http.StatusOK, http.StatusAccepted).
Expand All @@ -1579,6 +1621,43 @@ func (c *Client) GetBatchExportStatus(batchID, exportID string) (*BatchExportRes
return &response, nil
}

// ListBatchExportJobs returns every export job requested for a batch, most
// recent first. Manually deleted jobs are excluded.
func (c *Client) ListBatchExportJobs(batchID string) ([]ExportJobListItem, error) {
var response []ExportJobListItem

url := fmt.Sprintf("/api/v1/data-collection/batches/%s/export", batchID)
httpResponse, err := c.Execute(http.MethodGet, url, nil, &response)
if err != nil {
return nil, fmt.Errorf("unable to fulfil request %s: %s", url, err)
}

if httpResponse.StatusCode != http.StatusOK {
body, _ := io.ReadAll(httpResponse.Body)
return nil, fmt.Errorf("unexpected status code %d: %s", httpResponse.StatusCode, string(body))
}

return response, nil
}

// DeleteBatchExport permanently deletes an export job and, if it completed,
// its ZIP archive. An export that is still generating cannot be deleted (the
// API returns 409 Conflict).
func (c *Client) DeleteBatchExport(batchID, exportID string) error {
url := fmt.Sprintf("/api/v1/data-collection/batches/%s/export/%s", batchID, exportID)
httpResponse, err := c.Execute(http.MethodDelete, url, nil, nil)
if err != nil {
return fmt.Errorf("unable to fulfil request %s: %s", url, err)
}

if httpResponse.StatusCode != http.StatusNoContent {
body, _ := io.ReadAll(httpResponse.Body)
return fmt.Errorf("unexpected status code %d: %s", httpResponse.StatusCode, string(body))
}

return nil
}

// SyncAITaskBuilderBatch starts an async sync job that extends a batch with tasks
// created from datapoints appended to its dataset since setup or the last sync.
// Returns the created job (status "queued") including its sync_id.
Expand Down
61 changes: 61 additions & 0 deletions client/responses.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
package client

import (
"net/url"

"github.com/prolific-oss/cli/model"
)

Expand Down Expand Up @@ -465,6 +467,65 @@ type BatchExportResponse struct {
ExpiresAt string `json:"expires_at,omitempty"`
}

// ExportFilter narrows a batch/collection export to a subset of responses.
// All fields are optional; when every field is empty the export is
// unfiltered (a "full" export of every response).
type ExportFilter struct {
// StudyID restricts the export to responses submitted under this
// Prolific Study ID.
StudyID string
// From restricts the export to responses created on or after this ISO
// 8601 datetime (inclusive).
From string
// To restricts the export to responses created before this ISO 8601
// datetime (exclusive).
To string
}

// IsEmpty reports whether the filter has no fields set, i.e. it requests an
// unfiltered ("full") export.
func (f ExportFilter) IsEmpty() bool {
return f.StudyID == "" && f.From == "" && f.To == ""
}

// query encodes the filter as a URL query string. It returns an empty
// string when the filter is empty, so callers can append it unconditionally.
func (f ExportFilter) query() string {
if f.IsEmpty() {
return ""
}

values := url.Values{}
if f.StudyID != "" {
values.Set("study_id", f.StudyID)
}
if f.From != "" {
values.Set("from", f.From)
}
if f.To != "" {
values.Set("to", f.To)
}

return "?" + values.Encode()
}

// ExportJobFilter describes the filter (if any) an export job was requested
// with, as returned by the list-export-jobs endpoints.
type ExportJobFilter struct {
StudyID string `json:"study_id,omitempty"`
From string `json:"from,omitempty"`
To string `json:"to,omitempty"`
}

// ExportJobListItem summarizes one export job, as returned by the
// list-export-jobs endpoints for batches and collections.
type ExportJobListItem struct {
ExportID string `json:"export_id"`
Filter *ExportJobFilter `json:"filter"`
Status string `json:"status"`
CreatedAt string `json:"created_at"`
}

// AITaskBuilderBatchSyncResponse is the response for both starting a batch sync
// and polling its status. Status is one of "queued", "processing", "complete",
// or "failed". The outcome counts are populated on completion; Reason is set on
Expand Down
63 changes: 52 additions & 11 deletions cmd/aitaskbuilder/batch_export.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,8 +34,11 @@ var batchExportDownloadClient = http.DefaultClient

// BatchExportOptions holds the options for the batch export command.
type BatchExportOptions struct {
Args []string
Output string
Args []string
Output string
StudyID string
From string
To string
}

// NewBatchExportCommand creates a new `aitaskbuilder batch export` command to
Expand All @@ -59,7 +62,14 @@ the resulting ZIP file. The archive contains:
- files/ — participant-uploaded files (if any)

The export is generated asynchronously. This command will poll until the
archive is ready and then download it automatically.`,
archive is ready and then download it automatically.

Use --study-id, --from, and/or --to to narrow the export to a subset of
responses. Combining them applies all filters together (AND): --study-id
restricts the export to a specific Study, while --from/--to restrict it to
responses whose created_at falls within an ISO 8601 datetime range (--from
is inclusive, --to is exclusive). Omit all three for a full, unfiltered
export.`,
Example: `
Export a batch to the default filename (<batch-id>-export-<timestamp>.zip):

Expand All @@ -68,6 +78,14 @@ $ prolific aitaskbuilder batch export 5f8e3c2a-1d4b-4e6f-9a7c-2b0d8f3e1c5a
Export to a custom output path:

$ prolific aitaskbuilder batch export 5f8e3c2a-1d4b-4e6f-9a7c-2b0d8f3e1c5a --output /tmp/my-export.zip

Export only responses for a specific study:

$ prolific aitaskbuilder batch export 5f8e3c2a-1d4b-4e6f-9a7c-2b0d8f3e1c5a --study-id 60d3b2f1a2b3c4d5e6f7a8b9

Export only responses created in a date range:

$ prolific aitaskbuilder batch export 5f8e3c2a-1d4b-4e6f-9a7c-2b0d8f3e1c5a --from 2024-01-01T00:00:00Z --to 2024-02-01T00:00:00Z
`,
RunE: func(cmd *cobra.Command, args []string) error {
opts.Args = args
Expand All @@ -85,6 +103,15 @@ $ prolific aitaskbuilder batch export 5f8e3c2a-1d4b-4e6f-9a7c-2b0d8f3e1c5a --out
}

cmd.Flags().StringVarP(&opts.Output, "output", "o", "", "Output file path (default: <batch-id>-export-<timestamp>.zip)")
cmd.Flags().StringVar(&opts.StudyID, "study-id", "", "Restrict the export to responses submitted under this Study ID")
cmd.Flags().StringVar(&opts.From, "from", "", "Restrict the export to responses created on or after this ISO 8601 datetime (inclusive)")
cmd.Flags().StringVar(&opts.To, "to", "", "Restrict the export to responses created before this ISO 8601 datetime (exclusive)")

cmd.AddCommand(
NewBatchExportListCommand(c, w),
NewBatchExportDeleteCommand(c, w),
NewBatchExportDownloadCommand(c, w),
)

return cmd
}
Expand All @@ -95,7 +122,8 @@ func exportBatch(c client.API, opts BatchExportOptions, w io.Writer) error {
fmt.Fprintf(w, "Requesting export for batch %s...\n", batchID)

// Step 1: POST to initiate the export job.
initResult, err := c.InitiateBatchExport(batchID)
filter := client.ExportFilter{StudyID: opts.StudyID, From: opts.From, To: opts.To}
initResult, err := c.InitiateBatchExport(batchID, filter)
if err != nil {
return fmt.Errorf("error requesting export: %s", err.Error())
}
Expand All @@ -109,33 +137,46 @@ func exportBatch(c client.API, opts BatchExportOptions, w io.Writer) error {
return fmt.Errorf("unexpected export status %q for batch %s", initResult.Status, batchID)
}

exportID := initResult.ExportID

// Step 2: Poll GET until complete or failed.
url, err := pollBatchExportUntilDone(c, batchID, initResult.ExportID, w)
if err != nil {
return err
}

return batchDownloadExport(url, opts.Output, w)
}

// pollBatchExportUntilDone polls GetBatchExportStatus for batchID/exportID
// until the job reaches "complete" or "failed" (or the poll deadline is
// exceeded), printing a "." to w for each poll. It returns the presigned
// download URL on success. Shared by exportBatch (which polls a job it just
// requested) and downloadBatchExport (which polls a pre-existing job that
// was still generating).
func pollBatchExportUntilDone(c client.API, batchID, exportID string, w io.Writer) (string, error) {
deadline := time.Now().Add(batchExportTimeout)

for {
if time.Now().After(deadline) {
return fmt.Errorf("export timed out after 10 minutes for batch %s", batchID)
return "", fmt.Errorf("export timed out after 10 minutes for batch %s", batchID)
}

fmt.Fprint(w, ".")
batchExportPollSleep(batchExportPollInterval)

pollResult, err := c.GetBatchExportStatus(batchID, exportID)
if err != nil {
return fmt.Errorf("error polling export status: %s", err.Error())
return "", fmt.Errorf("error polling export status: %s", err.Error())
}

switch pollResult.Status {
case batchExportStatusComplete:
return batchDownloadExport(pollResult.URL, opts.Output, w)
return pollResult.URL, nil
case batchExportStatusFailed:
return fmt.Errorf("export generation failed for batch %s", batchID)
return "", fmt.Errorf("export generation failed for batch %s", batchID)
case batchExportStatusGenerating:
// continue polling
default:
return fmt.Errorf("unexpected export status %q for batch %s", pollResult.Status, batchID)
return "", fmt.Errorf("unexpected export status %q for batch %s", pollResult.Status, batchID)
}
}
}
Expand Down
Loading
Loading