diff --git a/README.md b/README.md index d58575a..1a401a2 100644 --- a/README.md +++ b/README.md @@ -64,7 +64,7 @@ client := loops.NewClient("YOUR_API_KEY", - Components — `GetComponent`, `ListComponents` - Themes — `GetTheme`, `ListThemes` - Uploads — `Upload`, `CreateUpload`, `CompleteUpload` -- Workflows — `ListWorkflows`, `GetWorkflow`, `GetWorkflowNode` +- Workflows — `ListWorkflows`, `GetWorkflow`, `GetWorkflowNode`, `CreateWorkflow`, `UpdateWorkflow`, `ChangeWorkflowMailingList`, `CreateWorkflowNode`, `UpdateWorkflowNode`, `AddWorkflowBranch`, `DeleteWorkflowNode`, `DeleteWorkflowNodeRecursive` Full reference: [pkg.go.dev/github.com/loops-so/loops-go](https://pkg.go.dev/github.com/loops-so/loops-go). diff --git a/workflows.go b/workflows.go index 17952d2..26c98d1 100644 --- a/workflows.go +++ b/workflows.go @@ -1,6 +1,7 @@ package loops import ( + "bytes" "encoding/json" "fmt" "net/http" @@ -16,17 +17,30 @@ type WorkflowSummary struct { } // SimplifiedWorkflow is returned by [Client.GetWorkflow]. Nodes is keyed by -// node ID; entry shapes are discriminated by their TypeName. +// node ID; entry shapes are discriminated by their TypeName. Status is one of +// the WorkflowStatus* constants. WorkflowRevisionID is the current revision +// token; pass it back as the expectedRevisionId on the next mutation (it may +// be nil for workflows without a revision token yet). type SimplifiedWorkflow struct { - ID string `json:"id"` - Name string `json:"name,omitempty"` - Description string `json:"description,omitempty"` - Emoji string `json:"emoji,omitempty"` - MailingListID *string `json:"mailingListId"` - RootNodeID *string `json:"rootNodeId"` - Nodes map[string]SimplifiedWorkflowNode `json:"nodes"` + ID string `json:"id"` + Status string `json:"status,omitempty"` + WorkflowRevisionID *string `json:"workflowRevisionId"` + Name string `json:"name,omitempty"` + Description string `json:"description,omitempty"` + Emoji string `json:"emoji,omitempty"` + MailingListID *string `json:"mailingListId"` + RootNodeID *string `json:"rootNodeId"` + Nodes map[string]SimplifiedWorkflowNode `json:"nodes"` } +// Workflow status values, used in [SimplifiedWorkflow]. +const ( + WorkflowStatusDraft = "Draft" + WorkflowStatusSending = "Sending" + WorkflowStatusPaused = "Paused" + WorkflowStatusPausedAndQueueing = "PausedAndQueueing" +) + // Workflow node typeName discriminator values, used in [WorkflowNode] and // [SimplifiedWorkflowNode]. const ( @@ -651,6 +665,29 @@ func marshalDiscriminated(typeName string, inner any) ([]byte, error) { return json.Marshal(fields) } +// mergeMarshal marshals inner (which has its own MarshalJSON) into a JSON +// object and overlays the extra fields. Used by the revision-bearing response +// types, whose promoted embedded MarshalJSON would otherwise drop their own +// fields. +func mergeMarshal(inner any, extra map[string]any) ([]byte, error) { + raw, err := json.Marshal(inner) + if err != nil { + return nil, err + } + var fields map[string]json.RawMessage + if err := json.Unmarshal(raw, &fields); err != nil { + return nil, err + } + for k, v := range extra { + b, err := json.Marshal(v) + if err != nil { + return nil, err + } + fields[k] = b + } + return json.Marshal(fields) +} + // ListWorkflows returns a single page of workflow summaries along with // pagination information. To iterate every page, use [Paginate]. func (c *Client) ListWorkflows(params PaginationParams) ([]WorkflowSummary, *Pagination, error) { @@ -718,8 +755,44 @@ func (c *Client) GetWorkflow(id string) (*SimplifiedWorkflow, error) { return &result, nil } -// GetWorkflowNode returns the detailed data for a single workflow node. -func (c *Client) GetWorkflowNode(workflowID, nodeID string) (*WorkflowNode, error) { +// WorkflowNodeWithRevision is returned by [Client.GetWorkflowNode]. It is the +// full [WorkflowNode] (its fields stay flat via the embedded type) plus the +// current workflow revision token. WorkflowRevisionID may be nil for workflows +// without a revision token yet. +type WorkflowNodeWithRevision struct { + WorkflowNode + WorkflowRevisionID *string `json:"workflowRevisionId"` +} + +// UnmarshalJSON decodes the node fields via the embedded [WorkflowNode] and +// then the revision token, since Go does not promote the embedded type's +// UnmarshalJSON when the outer struct has its own fields. +func (n *WorkflowNodeWithRevision) UnmarshalJSON(data []byte) error { + if err := n.WorkflowNode.UnmarshalJSON(data); err != nil { + return err + } + var rev struct { + WorkflowRevisionID *string `json:"workflowRevisionId"` + } + if err := json.Unmarshal(data, &rev); err != nil { + return err + } + n.WorkflowRevisionID = rev.WorkflowRevisionID + return nil +} + +// MarshalJSON encodes the node fields via the embedded [WorkflowNode] and adds +// the revision token, which the promoted embedded MarshalJSON would otherwise +// drop. +func (n WorkflowNodeWithRevision) MarshalJSON() ([]byte, error) { + return mergeMarshal(n.WorkflowNode, map[string]any{"workflowRevisionId": n.WorkflowRevisionID}) +} + +// GetWorkflowNode returns the detailed data for a single workflow node along +// with the current workflow revision token. The result is a +// [WorkflowNodeWithRevision]; the node fields remain accessible via the +// embedded [WorkflowNode]. +func (c *Client) GetWorkflowNode(workflowID, nodeID string) (*WorkflowNodeWithRevision, error) { req, err := c.newRequest(http.MethodGet, "/workflows/"+workflowID+"/nodes/"+nodeID, nil) if err != nil { return nil, err @@ -735,10 +808,860 @@ func (c *Client) GetWorkflowNode(workflowID, nodeID string) (*WorkflowNode, erro return nil, errorFromResponse(resp) } - var result WorkflowNode + var result WorkflowNodeWithRevision if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { return nil, fmt.Errorf("failed to decode response: %w", err) } return &result, nil } + +// WorkflowMutationNode is the detailed node returned by workflow mutations +// (create, update, add-branch). Unlike [WorkflowNode] it omits workflowId and +// its variant fields differ. Exactly one variant pointer is set; TypeName is +// the discriminator and matches the active variant (see the WorkflowNodeType* +// constants). +type WorkflowMutationNode struct { + TypeName string + SignupTrigger *SignupTriggerWorkflowMutationNode + EventTrigger *EventTriggerWorkflowMutationNode + ContactPropertyTrigger *ContactPropertyTriggerWorkflowMutationNode + AddToListTrigger *AddToListTriggerWorkflowMutationNode + BlankTrigger *BlankTriggerWorkflowMutationNode + AudienceFilter *AudienceFilterWorkflowMutationNode + TimerAction *TimerActionWorkflowMutationNode + SendEmailAction *SendEmailActionWorkflowMutationNode + ExitAction *ExitActionWorkflowMutationNode + BranchNode *BranchWorkflowMutationNode + ExperimentBranchNode *ExperimentBranchWorkflowMutationNode + VariantNode *VariantWorkflowMutationNode +} + +// SignupTriggerWorkflowMutationNode is the SignupTrigger variant of +// [WorkflowMutationNode]. +type SignupTriggerWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` +} + +// EventTriggerWorkflowMutationNode is the EventTrigger variant of +// [WorkflowMutationNode]. +type EventTriggerWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` + EventName string `json:"eventName,omitempty"` + EventProperties []WorkflowEventProperty `json:"eventProperties,omitempty"` + ReEligible bool `json:"reEligible"` +} + +// ContactPropertyTriggerWorkflowMutationNode is the ContactPropertyTrigger +// variant of [WorkflowMutationNode]. ContactPropertyQuery may be null. +type ContactPropertyTriggerWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` + ContactPropertyQuery *WorkflowContactPropertyQuery `json:"contactPropertyQuery"` + ReEligible bool `json:"reEligible"` +} + +// AddToListTriggerWorkflowMutationNode is the AddToListTrigger variant of +// [WorkflowMutationNode]. +type AddToListTriggerWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` + MailingListID *string `json:"mailingListId"` + ReEligible bool `json:"reEligible"` +} + +// BlankTriggerWorkflowMutationNode is the BlankTrigger variant of +// [WorkflowMutationNode]. +type BlankTriggerWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` +} + +// AudienceFilterWorkflowMutationNode is the AudienceFilter variant of +// [WorkflowMutationNode]. AudienceFilter reuses the package-level +// [AudienceFilter] type. +type AudienceFilterWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` + AudienceFilter *AudienceFilter `json:"audienceFilter,omitempty"` + AudienceSegmentID string `json:"audienceSegmentId,omitempty"` + AppliesDownstream bool `json:"appliesDownstream"` +} + +// TimerActionWorkflowMutationNode is the TimerAction variant of +// [WorkflowMutationNode]. +type TimerActionWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` + Amount float64 `json:"amount"` + Unit WorkflowTimerUnit `json:"unit"` +} + +// SendEmailActionWorkflowMutationNode is the SendEmailAction variant of +// [WorkflowMutationNode]. +type SendEmailActionWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` + EmailMessageID string `json:"emailMessageId"` + Subject string `json:"subject"` +} + +// ExitActionWorkflowMutationNode is the ExitAction variant of +// [WorkflowMutationNode]. +type ExitActionWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` +} + +// BranchWorkflowMutationNode is the BranchNode variant of +// [WorkflowMutationNode]. +type BranchWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` +} + +// ExperimentBranchWorkflowMutationNode is the ExperimentBranchNode variant of +// [WorkflowMutationNode]. +type ExperimentBranchWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` + SamplingRate float64 `json:"samplingRate"` +} + +// VariantWorkflowMutationNode is the VariantNode variant of +// [WorkflowMutationNode]. +type VariantWorkflowMutationNode struct { + ID string `json:"id"` + NextNodeIDs []string `json:"nextNodeIds"` + IsControl *bool `json:"isControl,omitempty"` +} + +// MarshalJSON encodes the active variant inline with a "typeName" +// discriminator. +func (n WorkflowMutationNode) MarshalJSON() ([]byte, error) { + inner, err := pickWorkflowMutationNode(n) + if err != nil { + return nil, err + } + return marshalDiscriminated(n.TypeName, inner) +} + +// UnmarshalJSON dispatches on "typeName" and decodes the matching variant. +func (n *WorkflowMutationNode) UnmarshalJSON(data []byte) error { + var head struct { + TypeName string `json:"typeName"` + } + if err := json.Unmarshal(data, &head); err != nil { + return err + } + n.TypeName = head.TypeName + switch head.TypeName { + case WorkflowNodeTypeSignupTrigger: + var v SignupTriggerWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.SignupTrigger = &v + case WorkflowNodeTypeEventTrigger: + var v EventTriggerWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.EventTrigger = &v + case WorkflowNodeTypeContactPropertyTrigger: + var v ContactPropertyTriggerWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.ContactPropertyTrigger = &v + case WorkflowNodeTypeAddToListTrigger: + var v AddToListTriggerWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.AddToListTrigger = &v + case WorkflowNodeTypeBlankTrigger: + var v BlankTriggerWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.BlankTrigger = &v + case WorkflowNodeTypeAudienceFilter: + var v AudienceFilterWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.AudienceFilter = &v + case WorkflowNodeTypeTimerAction: + var v TimerActionWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.TimerAction = &v + case WorkflowNodeTypeSendEmailAction: + var v SendEmailActionWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.SendEmailAction = &v + case WorkflowNodeTypeExitAction: + var v ExitActionWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.ExitAction = &v + case WorkflowNodeTypeBranchNode: + var v BranchWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.BranchNode = &v + case WorkflowNodeTypeExperimentBranchNode: + var v ExperimentBranchWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.ExperimentBranchNode = &v + case WorkflowNodeTypeVariantNode: + var v VariantWorkflowMutationNode + if err := json.Unmarshal(data, &v); err != nil { + return err + } + n.VariantNode = &v + default: + return fmt.Errorf("workflow node: unknown typeName %q", head.TypeName) + } + return nil +} + +func pickWorkflowMutationNode(n WorkflowMutationNode) (any, error) { + switch n.TypeName { + case WorkflowNodeTypeSignupTrigger: + return n.SignupTrigger, nil + case WorkflowNodeTypeEventTrigger: + return n.EventTrigger, nil + case WorkflowNodeTypeContactPropertyTrigger: + return n.ContactPropertyTrigger, nil + case WorkflowNodeTypeAddToListTrigger: + return n.AddToListTrigger, nil + case WorkflowNodeTypeBlankTrigger: + return n.BlankTrigger, nil + case WorkflowNodeTypeAudienceFilter: + return n.AudienceFilter, nil + case WorkflowNodeTypeTimerAction: + return n.TimerAction, nil + case WorkflowNodeTypeSendEmailAction: + return n.SendEmailAction, nil + case WorkflowNodeTypeExitAction: + return n.ExitAction, nil + case WorkflowNodeTypeBranchNode: + return n.BranchNode, nil + case WorkflowNodeTypeExperimentBranchNode: + return n.ExperimentBranchNode, nil + case WorkflowNodeTypeVariantNode: + return n.VariantNode, nil + } + return nil, fmt.Errorf("workflow node: unknown typeName %q", n.TypeName) +} + +// WorkflowMutationNodeWithRevision is a [WorkflowMutationNode] returned by a +// mutation (its fields stay flat via the embedded type) plus the current +// workflow revision token. +type WorkflowMutationNodeWithRevision struct { + WorkflowMutationNode + WorkflowRevisionID string `json:"workflowRevisionId"` +} + +// UnmarshalJSON decodes the node fields via the embedded +// [WorkflowMutationNode] and then the revision token. +func (n *WorkflowMutationNodeWithRevision) UnmarshalJSON(data []byte) error { + if err := n.WorkflowMutationNode.UnmarshalJSON(data); err != nil { + return err + } + var rev struct { + WorkflowRevisionID string `json:"workflowRevisionId"` + } + if err := json.Unmarshal(data, &rev); err != nil { + return err + } + n.WorkflowRevisionID = rev.WorkflowRevisionID + return nil +} + +// MarshalJSON encodes the node fields via the embedded [WorkflowMutationNode] +// and adds the revision token. +func (n WorkflowMutationNodeWithRevision) MarshalJSON() ([]byte, error) { + return mergeMarshal(n.WorkflowMutationNode, map[string]any{"workflowRevisionId": n.WorkflowRevisionID}) +} + +// CreatedWorkflowNode is the node returned by [Client.CreateWorkflowNode]. It +// is a [WorkflowMutationNodeWithRevision] (fields stay flat via the embedded +// type) plus any default child nodes created alongside the requested node — +// a BranchNode returns two AudienceFilter children; an ExperimentBranchNode +// returns two variant children and one control variant. +type CreatedWorkflowNode struct { + WorkflowMutationNodeWithRevision + CreatedChildNodes []WorkflowMutationNode `json:"createdChildNodes,omitempty"` +} + +// UnmarshalJSON decodes the node and revision via the embedded type and then +// the createdChildNodes list. +func (n *CreatedWorkflowNode) UnmarshalJSON(data []byte) error { + if err := n.WorkflowMutationNodeWithRevision.UnmarshalJSON(data); err != nil { + return err + } + var extra struct { + CreatedChildNodes []WorkflowMutationNode `json:"createdChildNodes"` + } + if err := json.Unmarshal(data, &extra); err != nil { + return err + } + n.CreatedChildNodes = extra.CreatedChildNodes + return nil +} + +// MarshalJSON encodes the embedded [WorkflowMutationNodeWithRevision] and adds +// the createdChildNodes list when present. +func (n CreatedWorkflowNode) MarshalJSON() ([]byte, error) { + extra := map[string]any{} + if len(n.CreatedChildNodes) > 0 { + extra["createdChildNodes"] = n.CreatedChildNodes + } + return mergeMarshal(n.WorkflowMutationNodeWithRevision, extra) +} + +// CreateWorkflowRequest is the request body for [Client.CreateWorkflow]. +// MailingListID is nullable. +type CreateWorkflowRequest struct { + Name string `json:"name"` + Description string `json:"description,omitempty"` + MailingListID *string `json:"mailingListId"` +} + +// CreateWorkflow creates a new workflow and returns it. +func (c *Client) CreateWorkflow(req CreateWorkflowRequest) (*SimplifiedWorkflow, error) { + b, err := json.Marshal(req) + if err != nil { + return nil, fmt.Errorf("failed to encode request: %w", err) + } + + httpReq, err := c.newRequest(http.MethodPost, "/workflows", bytes.NewReader(b)) + if err != nil { + return nil, err + } + + resp, err := c.do(httpReq) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, errorFromResponse(resp) + } + + var result SimplifiedWorkflow + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return nil, fmt.Errorf("failed to decode response: %w", err) + } + + return &result, nil +} + +// UpdateWorkflowPropertiesRequest is the request body for +// [Client.UpdateWorkflow]. ExpectedRevisionID is required for optimistic +// concurrency and is always sent (as JSON null when nil); a stale value +// returns a 409. At least one of Name or Description should be set. +type UpdateWorkflowPropertiesRequest struct { + ExpectedRevisionID *string + Name string + Description string +} + +// UpdateWorkflow updates the name and/or description of the workflow +// identified by id and returns its new state. +func (c *Client) UpdateWorkflow(id string, req UpdateWorkflowPropertiesRequest) (*SimplifiedWorkflow, error) { + body := map[string]any{ + "expectedRevisionId": req.ExpectedRevisionID, + } + if req.Name != "" { + body["name"] = req.Name + } + if req.Description != "" { + body["description"] = req.Description + } + + b, err := json.Marshal(body) + if err != nil { + return nil, fmt.Errorf("failed to encode request: %w", err) + } + + httpReq, err := c.newRequest(http.MethodPost, "/workflows/"+id, bytes.NewReader(b)) + if err != nil { + return nil, err + } + + resp, err := c.do(httpReq) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, errorFromResponse(resp) + } + + var result SimplifiedWorkflow + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return nil, fmt.Errorf("failed to decode response: %w", err) + } + + return &result, nil +} + +// Queued-contact policies controlling how a mutation treats queued contacts +// that would be excluded by the change. +const ( + // WorkflowQueuedContactPolicyFail returns queued-contact impact instead of + // mutating. It is the default when omitted. + WorkflowQueuedContactPolicyFail = "fail" + // WorkflowQueuedContactPolicyDiscard confirms that matching queued + // contacts should be discarded. + WorkflowQueuedContactPolicyDiscard = "discard" +) + +// Workflow mailing-list and node-deletion status values, returned in +// [ChangeWorkflowMailingListResponse] and [DeleteWorkflowNodeResponse]. +const ( + WorkflowMutationStatusDryRun = "dryRun" + WorkflowMutationStatusQueuedContactsFound = "queuedContactsFound" + WorkflowMutationStatusUpdated = "updated" + WorkflowMutationStatusDeleted = "deleted" +) + +// ChangeWorkflowMailingListRequest is the request body for +// [Client.ChangeWorkflowMailingList]. ExpectedRevisionID and MailingListID are +// both required and nullable and always sent (as JSON null when nil). +type ChangeWorkflowMailingListRequest struct { + ExpectedRevisionID *string + MailingListID *string + DryRun bool + QueuedContactPolicy string +} + +// ChangeWorkflowMailingListResponse is returned by +// [Client.ChangeWorkflowMailingList]. Status is one of the +// WorkflowMutationStatus* values. WorkflowRevisionID is set only when Status +// is "updated". +type ChangeWorkflowMailingListResponse struct { + Status string `json:"status"` + MailingListID *string `json:"mailingListId"` + WorkflowRevisionID *string `json:"workflowRevisionId,omitempty"` + QueuedContactCount float64 `json:"queuedContactCount"` +} + +// ChangeWorkflowMailingList changes the mailing list of the workflow +// identified by id. With the default "fail" policy (or DryRun), a response +// whose Status is not "updated" reports the queued-contact impact instead of +// applying the change. +func (c *Client) ChangeWorkflowMailingList(id string, req ChangeWorkflowMailingListRequest) (*ChangeWorkflowMailingListResponse, error) { + body := map[string]any{ + "expectedRevisionId": req.ExpectedRevisionID, + "mailingListId": req.MailingListID, + } + if req.DryRun { + body["dryRun"] = req.DryRun + } + if req.QueuedContactPolicy != "" { + body["queuedContactPolicy"] = req.QueuedContactPolicy + } + + b, err := json.Marshal(body) + if err != nil { + return nil, fmt.Errorf("failed to encode request: %w", err) + } + + httpReq, err := c.newRequest(http.MethodPost, "/workflows/"+id+"/mailing-list", bytes.NewReader(b)) + if err != nil { + return nil, err + } + + resp, err := c.do(httpReq) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, errorFromResponse(resp) + } + + var result ChangeWorkflowMailingListResponse + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return nil, fmt.Errorf("failed to decode response: %w", err) + } + + return &result, nil +} + +// Workflow node insert modes for [CreateWorkflowNodeRequest]. +const ( + WorkflowInsertModeBetween = "between" + WorkflowInsertModeBefore = "before" +) + +// Node types that can be created via [Client.CreateWorkflowNode]. Trigger and +// ExitAction nodes cannot be created. +const ( + CreateWorkflowNodeTypeAudienceFilter = WorkflowNodeTypeAudienceFilter + CreateWorkflowNodeTypeBranchNode = WorkflowNodeTypeBranchNode + CreateWorkflowNodeTypeExperimentBranchNode = WorkflowNodeTypeExperimentBranchNode + CreateWorkflowNodeTypeTimerAction = WorkflowNodeTypeTimerAction + CreateWorkflowNodeTypeSendEmailAction = WorkflowNodeTypeSendEmailAction + CreateWorkflowNodeTypeVariantNode = WorkflowNodeTypeVariantNode +) + +// CreateWorkflowNodeRequest is the request body for [Client.CreateWorkflowNode]. +// InsertMode selects the placement: +// - "between" inserts between FromNodeID and ToNodeID (FromNodeID must +// currently point to ToNodeID); both are required. +// - "before" inserts before BeforeNodeID, which must have an incoming parent +// and cannot be a trigger node. +// +// ExpectedRevisionID is required and always sent (as JSON null when nil). +// NodeTypeName is one of the CreateWorkflowNodeType* constants. +type CreateWorkflowNodeRequest struct { + ExpectedRevisionID *string + InsertMode string + NodeTypeName string + FromNodeID string + ToNodeID string + BeforeNodeID string +} + +// CreateWorkflowNodeResponse is returned by [Client.CreateWorkflowNode]. +type CreateWorkflowNodeResponse struct { + Node CreatedWorkflowNode `json:"node"` + Workflow SimplifiedWorkflow `json:"workflow"` +} + +// CreateWorkflowNode creates a node in the workflow identified by id and +// returns the created node together with the updated workflow. Configure the +// node afterwards with [Client.UpdateWorkflowNode]. +func (c *Client) CreateWorkflowNode(id string, req CreateWorkflowNodeRequest) (*CreateWorkflowNodeResponse, error) { + body := map[string]any{ + "expectedRevisionId": req.ExpectedRevisionID, + "insertMode": req.InsertMode, + "nodeTypeName": req.NodeTypeName, + } + switch req.InsertMode { + case WorkflowInsertModeBetween: + body["fromNodeId"] = req.FromNodeID + body["toNodeId"] = req.ToNodeID + case WorkflowInsertModeBefore: + body["beforeNodeId"] = req.BeforeNodeID + } + + b, err := json.Marshal(body) + if err != nil { + return nil, fmt.Errorf("failed to encode request: %w", err) + } + + httpReq, err := c.newRequest(http.MethodPost, "/workflows/"+id+"/nodes", bytes.NewReader(b)) + if err != nil { + return nil, err + } + + resp, err := c.do(httpReq) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, errorFromResponse(resp) + } + + var result CreateWorkflowNodeResponse + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return nil, fmt.Errorf("failed to decode response: %w", err) + } + + return &result, nil +} + +// UpdateWorkflowNodePayload is the node-type-specific configuration applied by +// [Client.UpdateWorkflowNode]. Exactly one variant pointer is set; TypeName is +// the discriminator and matches the active variant (see the WorkflowNodeType* +// constants). Trigger payloads emit a "typeName" field (which may change one +// trigger type to another); the AudienceFilter, TimerAction, ExperimentBranch, +// and Variant config payloads do not. +type UpdateWorkflowNodePayload struct { + TypeName string + SignupTrigger *WorkflowSignupTriggerPayload + EventTrigger *WorkflowEventTriggerPayload + ContactPropertyTrigger *WorkflowContactPropertyTriggerPayload + AddToListTrigger *WorkflowAddToListTriggerPayload + AudienceFilter *WorkflowAudienceFilterPayload + TimerAction *WorkflowTimerActionPayload + ExperimentBranch *WorkflowExperimentBranchPayload + Variant *WorkflowVariantPayload +} + +// WorkflowSignupTriggerPayload changes an existing trigger node to a signup +// trigger. +type WorkflowSignupTriggerPayload struct{} + +// WorkflowEventTriggerPayload updates an event trigger, or changes a trigger +// node to an event trigger. Set EventPatternID or EventName (not both); either +// may be explicit JSON null to clear the relationship. +type WorkflowEventTriggerPayload struct { + EventPatternID *string `json:"eventPatternId,omitempty"` + EventName *string `json:"eventName,omitempty"` + ReEligible *bool `json:"reEligible,omitempty"` +} + +// WorkflowContactPropertyTriggerPayload updates a contact-property trigger, or +// changes a trigger node to one. ContactPropertyQuery reuses the package-level +// [WorkflowContactPropertyQuery] type. +type WorkflowContactPropertyTriggerPayload struct { + ContactPropertyQuery *WorkflowContactPropertyQuery `json:"contactPropertyQuery,omitempty"` + ReEligible *bool `json:"reEligible,omitempty"` +} + +// WorkflowAddToListTriggerPayload updates an add-to-list trigger, or changes a +// trigger node to one. +type WorkflowAddToListTriggerPayload struct { + ReEligible *bool `json:"reEligible,omitempty"` +} + +// WorkflowAudienceFilterPayload configures an audience filter node. +// AudienceFilter reuses the package-level [AudienceFilter] type. +type WorkflowAudienceFilterPayload struct { + AudienceSegmentID *string `json:"audienceSegmentId,omitempty"` + AudienceFilter *AudienceFilter `json:"audienceFilter,omitempty"` + AppliesDownstream *bool `json:"appliesDownstream,omitempty"` +} + +// WorkflowTimerActionPayload configures a timer action node. +type WorkflowTimerActionPayload struct { + Amount *float64 `json:"amount,omitempty"` + Unit WorkflowTimerUnit `json:"unit,omitempty"` +} + +// WorkflowExperimentBranchPayload configures an experiment branch node. +// SamplingRate is a percentage between 0 and 100. +type WorkflowExperimentBranchPayload struct { + SamplingRate *float64 `json:"samplingRate,omitempty"` +} + +// WorkflowVariantPayload configures a variant node. +type WorkflowVariantPayload struct { + IsControl *bool `json:"isControl,omitempty"` +} + +// MarshalJSON encodes the active payload variant. Trigger variants include a +// "typeName" discriminator; config variants are emitted as-is. +func (p UpdateWorkflowNodePayload) MarshalJSON() ([]byte, error) { + switch p.TypeName { + case WorkflowNodeTypeSignupTrigger: + return marshalDiscriminated(p.TypeName, p.SignupTrigger) + case WorkflowNodeTypeEventTrigger: + return marshalDiscriminated(p.TypeName, p.EventTrigger) + case WorkflowNodeTypeContactPropertyTrigger: + return marshalDiscriminated(p.TypeName, p.ContactPropertyTrigger) + case WorkflowNodeTypeAddToListTrigger: + return marshalDiscriminated(p.TypeName, p.AddToListTrigger) + case WorkflowNodeTypeAudienceFilter: + if p.AudienceFilter == nil { + return nil, fmt.Errorf("workflow node payload: %s variant is nil", p.TypeName) + } + return json.Marshal(p.AudienceFilter) + case WorkflowNodeTypeTimerAction: + if p.TimerAction == nil { + return nil, fmt.Errorf("workflow node payload: %s variant is nil", p.TypeName) + } + return json.Marshal(p.TimerAction) + case WorkflowNodeTypeExperimentBranchNode: + if p.ExperimentBranch == nil { + return nil, fmt.Errorf("workflow node payload: %s variant is nil", p.TypeName) + } + return json.Marshal(p.ExperimentBranch) + case WorkflowNodeTypeVariantNode: + if p.Variant == nil { + return nil, fmt.Errorf("workflow node payload: %s variant is nil", p.TypeName) + } + return json.Marshal(p.Variant) + } + return nil, fmt.Errorf("workflow node payload: unknown typeName %q", p.TypeName) +} + +// UpdateWorkflowNodeRequest is the request body for [Client.UpdateWorkflowNode]. +// ExpectedRevisionID is required and always sent (as JSON null when nil). +type UpdateWorkflowNodeRequest struct { + ExpectedRevisionID *string + Payload UpdateWorkflowNodePayload +} + +// UpdateWorkflowNodeResponse is the updated node returned by +// [Client.UpdateWorkflowNode], with the latest workflow revision token. +type UpdateWorkflowNodeResponse = WorkflowMutationNodeWithRevision + +// UpdateWorkflowNode applies req.Payload to the node identified by +// workflowID/nodeID and returns the updated node. +func (c *Client) UpdateWorkflowNode(workflowID, nodeID string, req UpdateWorkflowNodeRequest) (*UpdateWorkflowNodeResponse, error) { + body := map[string]any{ + "expectedRevisionId": req.ExpectedRevisionID, + "payload": req.Payload, + } + + b, err := json.Marshal(body) + if err != nil { + return nil, fmt.Errorf("failed to encode request: %w", err) + } + + httpReq, err := c.newRequest(http.MethodPost, "/workflows/"+workflowID+"/nodes/"+nodeID, bytes.NewReader(b)) + if err != nil { + return nil, err + } + + resp, err := c.do(httpReq) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, errorFromResponse(resp) + } + + var result UpdateWorkflowNodeResponse + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return nil, fmt.Errorf("failed to decode response: %w", err) + } + + return &result, nil +} + +// AddWorkflowBranchRequest is the request body for [Client.AddWorkflowBranch]. +// ExpectedRevisionID is required and always sent (as JSON null when nil). +type AddWorkflowBranchRequest struct { + ExpectedRevisionID *string +} + +// AddWorkflowBranchResponse is returned by [Client.AddWorkflowBranch]. +type AddWorkflowBranchResponse struct { + Node WorkflowMutationNodeWithRevision `json:"node"` + Workflow SimplifiedWorkflow `json:"workflow"` +} + +// AddWorkflowBranch adds a branch to the branch node identified by +// workflowID/nodeID and returns the updated node and workflow. +func (c *Client) AddWorkflowBranch(workflowID, nodeID string, req AddWorkflowBranchRequest) (*AddWorkflowBranchResponse, error) { + body := map[string]any{ + "expectedRevisionId": req.ExpectedRevisionID, + } + + b, err := json.Marshal(body) + if err != nil { + return nil, fmt.Errorf("failed to encode request: %w", err) + } + + httpReq, err := c.newRequest(http.MethodPost, "/workflows/"+workflowID+"/nodes/"+nodeID+"/add-branch", bytes.NewReader(b)) + if err != nil { + return nil, err + } + + resp, err := c.do(httpReq) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, errorFromResponse(resp) + } + + var result AddWorkflowBranchResponse + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return nil, fmt.Errorf("failed to decode response: %w", err) + } + + return &result, nil +} + +// DeleteWorkflowNodeRequest is the request body for [Client.DeleteWorkflowNode] +// and [Client.DeleteWorkflowNodeRecursive]. ExpectedRevisionID is required and +// always sent (as JSON null when nil). QueuedContactPolicy is one of the +// WorkflowQueuedContactPolicy* constants. +type DeleteWorkflowNodeRequest struct { + ExpectedRevisionID *string + DryRun bool + QueuedContactPolicy string +} + +// DeleteWorkflowNodeResponse is returned by [Client.DeleteWorkflowNode] and +// [Client.DeleteWorkflowNodeRecursive]. Status is one of the +// WorkflowMutationStatus* values. WorkflowRevisionID is set only when Status +// is "deleted". +type DeleteWorkflowNodeResponse struct { + Status string `json:"status"` + NodeIDs []string `json:"nodeIds"` + WorkflowRevisionID *string `json:"workflowRevisionId,omitempty"` + QueuedContactCount float64 `json:"queuedContactCount"` +} + +func (c *Client) deleteWorkflowNode(path string, req DeleteWorkflowNodeRequest) (*DeleteWorkflowNodeResponse, error) { + body := map[string]any{ + "expectedRevisionId": req.ExpectedRevisionID, + } + if req.DryRun { + body["dryRun"] = req.DryRun + } + if req.QueuedContactPolicy != "" { + body["queuedContactPolicy"] = req.QueuedContactPolicy + } + + b, err := json.Marshal(body) + if err != nil { + return nil, fmt.Errorf("failed to encode request: %w", err) + } + + httpReq, err := c.newRequest(http.MethodDelete, path, bytes.NewReader(b)) + if err != nil { + return nil, err + } + + resp, err := c.do(httpReq) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, errorFromResponse(resp) + } + + var result DeleteWorkflowNodeResponse + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return nil, fmt.Errorf("failed to decode response: %w", err) + } + + return &result, nil +} + +// DeleteWorkflowNode deletes the node identified by workflowID/nodeID. With the +// default "fail" policy (or DryRun), a response whose Status is not "deleted" +// reports the queued-contact impact instead of applying the deletion. +func (c *Client) DeleteWorkflowNode(workflowID, nodeID string, req DeleteWorkflowNodeRequest) (*DeleteWorkflowNodeResponse, error) { + return c.deleteWorkflowNode("/workflows/"+workflowID+"/nodes/"+nodeID, req) +} + +// DeleteWorkflowNodeRecursive deletes the node identified by workflowID/nodeID +// along with its downstream nodes. Semantics otherwise match +// [Client.DeleteWorkflowNode]. +func (c *Client) DeleteWorkflowNodeRecursive(workflowID, nodeID string, req DeleteWorkflowNodeRequest) (*DeleteWorkflowNodeResponse, error) { + return c.deleteWorkflowNode("/workflows/"+workflowID+"/nodes/"+nodeID+"/recursive", req) +} diff --git a/workflows_test.go b/workflows_test.go index cc7c74c..d88b6f7 100644 --- a/workflows_test.go +++ b/workflows_test.go @@ -3,6 +3,7 @@ package loops import ( "encoding/json" "errors" + "io" "net/http" "net/http/httptest" "strings" @@ -36,6 +37,8 @@ const listWorkflowsResponse = `{ const getSimplifiedWorkflowResponse = `{ "id": "wf_1", + "status": "Draft", + "workflowRevisionId": "rev_1", "name": "Onboarding", "description": "Welcome new signups", "emoji": "👋", @@ -189,6 +192,12 @@ func TestGetWorkflow(t *testing.T) { if result.ID != "wf_1" { t.Errorf("ID = %q, want wf_1", result.ID) } + if result.Status != WorkflowStatusDraft { + t.Errorf("Status = %q, want Draft", result.Status) + } + if result.WorkflowRevisionID == nil || *result.WorkflowRevisionID != "rev_1" { + t.Errorf("WorkflowRevisionID = %v, want rev_1", result.WorkflowRevisionID) + } if result.RootNodeID == nil || *result.RootNodeID != "node_root" { t.Errorf("RootNodeID = %v, want node_root", result.RootNodeID) } @@ -248,7 +257,8 @@ func TestGetWorkflowNode_EventTrigger(t *testing.T) { {"name": "plan", "type": "string"}, {"name": "seats", "type": "number"} ], - "reEligible": true + "reEligible": true, + "workflowRevisionId": "rev_1" }` server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if want := "/workflows/wf_1/nodes/node_evt"; r.URL.Path != want { @@ -264,6 +274,9 @@ func TestGetWorkflowNode_EventTrigger(t *testing.T) { if err != nil { t.Fatalf("unexpected error: %v", err) } + if node.WorkflowRevisionID == nil || *node.WorkflowRevisionID != "rev_1" { + t.Errorf("WorkflowRevisionID = %v, want rev_1", node.WorkflowRevisionID) + } if node.TypeName != WorkflowNodeTypeEventTrigger { t.Fatalf("TypeName = %q, want EventTrigger", node.TypeName) } @@ -488,6 +501,618 @@ func TestSimplifiedWorkflowNode_MarshalRoundTrip(t *testing.T) { } } +// decodeBody reads and JSON-decodes the request body of a test server handler. +func decodeBody(t *testing.T, r *http.Request) map[string]any { + t.Helper() + var got map[string]any + if err := json.NewDecoder(r.Body).Decode(&got); err != nil { + t.Fatalf("decode request body: %v", err) + } + return got +} + +func TestCreateWorkflow(t *testing.T) { + var gotPath, gotMethod string + var gotBody map[string]any + resp := `{ + "id": "wf_new", + "status": "Draft", + "workflowRevisionId": null, + "name": "New flow", + "mailingListId": "ml_1", + "rootNodeId": "node_root", + "nodes": {} + }` + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotPath = r.URL.Path + gotMethod = r.Method + gotBody = decodeBody(t, r) + w.WriteHeader(http.StatusOK) + w.Write([]byte(resp)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + wf, err := client.CreateWorkflow(CreateWorkflowRequest{ + Name: "New flow", + Description: "desc", + MailingListID: ptr("ml_1"), + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if gotMethod != http.MethodPost { + t.Errorf("method = %q, want POST", gotMethod) + } + if gotPath != "/workflows" { + t.Errorf("path = %q, want /workflows", gotPath) + } + if gotBody["name"] != "New flow" || gotBody["description"] != "desc" || gotBody["mailingListId"] != "ml_1" { + t.Errorf("body = %+v", gotBody) + } + if wf.ID != "wf_new" || wf.Status != WorkflowStatusDraft { + t.Errorf("wf = %+v", wf) + } + if wf.WorkflowRevisionID != nil { + t.Errorf("WorkflowRevisionID = %v, want nil", wf.WorkflowRevisionID) + } +} + +func TestUpdateWorkflow_NullRevision(t *testing.T) { + var gotBody map[string]any + var rawBody string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + buf := new(strings.Builder) + tee := io.TeeReader(r.Body, buf) + json.NewDecoder(tee).Decode(&gotBody) + rawBody = buf.String() + if want := "/workflows/wf_1"; r.URL.Path != want { + t.Errorf("path = %q, want %q", r.URL.Path, want) + } + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"id":"wf_1","status":"Draft","workflowRevisionId":"rev_2","name":"Renamed","mailingListId":null,"rootNodeId":"n","nodes":{}}`)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + wf, err := client.UpdateWorkflow("wf_1", UpdateWorkflowPropertiesRequest{ + ExpectedRevisionID: nil, + Name: "Renamed", + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !strings.Contains(rawBody, `"expectedRevisionId":null`) { + t.Errorf("body must send explicit null revision: %s", rawBody) + } + if _, ok := gotBody["expectedRevisionId"]; !ok { + t.Error("expectedRevisionId key missing from body") + } + if gotBody["name"] != "Renamed" { + t.Errorf("name = %v, want Renamed", gotBody["name"]) + } + if _, ok := gotBody["description"]; ok { + t.Errorf("description should be omitted, body = %+v", gotBody) + } + if wf.Name != "Renamed" || wf.WorkflowRevisionID == nil || *wf.WorkflowRevisionID != "rev_2" { + t.Errorf("wf = %+v", wf) + } +} + +func TestUpdateWorkflow_StaleRevision409(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusConflict) + w.Write([]byte(`{"message":"Workflow revision is stale."}`)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + _, err := client.UpdateWorkflow("wf_1", UpdateWorkflowPropertiesRequest{ + ExpectedRevisionID: ptr("old_rev"), + Name: "x", + }) + var apiErr *APIError + if !errors.As(err, &apiErr) { + t.Fatalf("expected *APIError, got %T: %v", err, err) + } + if apiErr.StatusCode != http.StatusConflict { + t.Errorf("StatusCode = %d, want 409", apiErr.StatusCode) + } + if apiErr.Message != "Workflow revision is stale." { + t.Errorf("Message = %q", apiErr.Message) + } +} + +func TestChangeWorkflowMailingList_Preview(t *testing.T) { + var gotBody map[string]any + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody = decodeBody(t, r) + if want := "/workflows/wf_1/mailing-list"; r.URL.Path != want { + t.Errorf("path = %q, want %q", r.URL.Path, want) + } + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"status":"queuedContactsFound","mailingListId":"ml_2","queuedContactCount":7}`)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + res, err := client.ChangeWorkflowMailingList("wf_1", ChangeWorkflowMailingListRequest{ + ExpectedRevisionID: ptr("rev_1"), + MailingListID: ptr("ml_2"), + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if gotBody["expectedRevisionId"] != "rev_1" || gotBody["mailingListId"] != "ml_2" { + t.Errorf("body = %+v", gotBody) + } + if _, ok := gotBody["dryRun"]; ok { + t.Errorf("dryRun should be omitted, body = %+v", gotBody) + } + if res.Status != WorkflowMutationStatusQueuedContactsFound { + t.Errorf("Status = %q, want queuedContactsFound", res.Status) + } + if res.QueuedContactCount != 7 { + t.Errorf("QueuedContactCount = %v, want 7", res.QueuedContactCount) + } + if res.WorkflowRevisionID != nil { + t.Errorf("WorkflowRevisionID = %v, want nil", res.WorkflowRevisionID) + } +} + +func TestChangeWorkflowMailingList_NullClear(t *testing.T) { + var rawBody string + var gotBody map[string]any + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + buf := new(strings.Builder) + tee := io.TeeReader(r.Body, buf) + json.NewDecoder(tee).Decode(&gotBody) + rawBody = buf.String() + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"status":"updated","mailingListId":null,"workflowRevisionId":"rev_3","queuedContactCount":0}`)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + res, err := client.ChangeWorkflowMailingList("wf_1", ChangeWorkflowMailingListRequest{ + ExpectedRevisionID: ptr("rev_2"), + MailingListID: nil, + DryRun: true, + QueuedContactPolicy: WorkflowQueuedContactPolicyDiscard, + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !strings.Contains(rawBody, `"mailingListId":null`) { + t.Errorf("body must send explicit null mailingListId: %s", rawBody) + } + if gotBody["dryRun"] != true || gotBody["queuedContactPolicy"] != "discard" { + t.Errorf("body = %+v", gotBody) + } + if res.Status != WorkflowMutationStatusUpdated { + t.Errorf("Status = %q, want updated", res.Status) + } + if res.WorkflowRevisionID == nil || *res.WorkflowRevisionID != "rev_3" { + t.Errorf("WorkflowRevisionID = %v, want rev_3", res.WorkflowRevisionID) + } + if res.MailingListID != nil { + t.Errorf("MailingListID = %v, want nil", res.MailingListID) + } +} + +func TestCreateWorkflowNode_Between(t *testing.T) { + var gotBody map[string]any + resp := `{ + "node": { + "id": "node_new", + "typeName": "BranchNode", + "nextNodeIds": ["c1", "c2"], + "workflowRevisionId": "rev_5", + "createdChildNodes": [ + {"id": "c1", "typeName": "AudienceFilter", "nextNodeIds": [], "appliesDownstream": false}, + {"id": "c2", "typeName": "AudienceFilter", "nextNodeIds": [], "appliesDownstream": false} + ] + }, + "workflow": {"id":"wf_1","status":"Draft","workflowRevisionId":"rev_5","mailingListId":null,"rootNodeId":"r","nodes":{}} + }` + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody = decodeBody(t, r) + if want := "/workflows/wf_1/nodes"; r.URL.Path != want { + t.Errorf("path = %q, want %q", r.URL.Path, want) + } + w.WriteHeader(http.StatusOK) + w.Write([]byte(resp)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + res, err := client.CreateWorkflowNode("wf_1", CreateWorkflowNodeRequest{ + ExpectedRevisionID: ptr("rev_4"), + InsertMode: WorkflowInsertModeBetween, + NodeTypeName: CreateWorkflowNodeTypeBranchNode, + FromNodeID: "node_a", + ToNodeID: "node_b", + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if gotBody["insertMode"] != "between" || gotBody["fromNodeId"] != "node_a" || gotBody["toNodeId"] != "node_b" { + t.Errorf("body = %+v", gotBody) + } + if gotBody["nodeTypeName"] != "BranchNode" { + t.Errorf("nodeTypeName = %v", gotBody["nodeTypeName"]) + } + if _, ok := gotBody["beforeNodeId"]; ok { + t.Errorf("beforeNodeId should be absent for between: %+v", gotBody) + } + if res.Node.TypeName != WorkflowNodeTypeBranchNode || res.Node.BranchNode == nil { + t.Errorf("node = %+v", res.Node) + } + if res.Node.WorkflowRevisionID != "rev_5" { + t.Errorf("node revision = %q, want rev_5", res.Node.WorkflowRevisionID) + } + if len(res.Node.CreatedChildNodes) != 2 { + t.Fatalf("len(CreatedChildNodes) = %d, want 2", len(res.Node.CreatedChildNodes)) + } + if res.Node.CreatedChildNodes[0].TypeName != WorkflowNodeTypeAudienceFilter || res.Node.CreatedChildNodes[0].AudienceFilter == nil { + t.Errorf("child[0] = %+v", res.Node.CreatedChildNodes[0]) + } + if res.Workflow.ID != "wf_1" { + t.Errorf("workflow.ID = %q", res.Workflow.ID) + } +} + +func TestCreateWorkflowNode_Before_NullRevision(t *testing.T) { + var rawBody string + var gotBody map[string]any + resp := `{ + "node": {"id":"node_new","typeName":"TimerAction","nextNodeIds":[],"amount":0,"unit":"m","workflowRevisionId":"rev_1"}, + "workflow": {"id":"wf_1","status":"Draft","workflowRevisionId":"rev_1","mailingListId":null,"rootNodeId":"r","nodes":{}} + }` + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + buf := new(strings.Builder) + tee := io.TeeReader(r.Body, buf) + json.NewDecoder(tee).Decode(&gotBody) + rawBody = buf.String() + w.WriteHeader(http.StatusOK) + w.Write([]byte(resp)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + res, err := client.CreateWorkflowNode("wf_1", CreateWorkflowNodeRequest{ + ExpectedRevisionID: nil, + InsertMode: WorkflowInsertModeBefore, + NodeTypeName: CreateWorkflowNodeTypeTimerAction, + BeforeNodeID: "node_target", + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !strings.Contains(rawBody, `"expectedRevisionId":null`) { + t.Errorf("body must send explicit null revision: %s", rawBody) + } + if gotBody["insertMode"] != "before" || gotBody["beforeNodeId"] != "node_target" { + t.Errorf("body = %+v", gotBody) + } + if _, ok := gotBody["fromNodeId"]; ok { + t.Errorf("fromNodeId should be absent for before: %+v", gotBody) + } + if res.Node.TypeName != WorkflowNodeTypeTimerAction || res.Node.TimerAction == nil { + t.Errorf("node = %+v", res.Node) + } + if res.Node.TimerAction.Unit != WorkflowTimerUnitMinutes { + t.Errorf("unit = %q, want m", res.Node.TimerAction.Unit) + } +} + +func TestUpdateWorkflowNode_Payloads(t *testing.T) { + tests := []struct { + name string + payload UpdateWorkflowNodePayload + // wantSubstr are JSON fragments that must appear in the marshaled payload. + wantSubstr []string + // notSubstr are fragments that must NOT appear. + notSubstr []string + }{ + { + name: "SignupTrigger has typeName", + payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeSignupTrigger, + SignupTrigger: &WorkflowSignupTriggerPayload{}, + }, + wantSubstr: []string{`"typeName":"SignupTrigger"`}, + }, + { + name: "EventTrigger with null eventName", + payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeEventTrigger, + EventTrigger: &WorkflowEventTriggerPayload{EventName: nil, EventPatternID: ptr("ep_1"), ReEligible: ptr(true)}, + }, + wantSubstr: []string{`"typeName":"EventTrigger"`, `"eventPatternId":"ep_1"`, `"reEligible":true`}, + }, + { + name: "ContactPropertyTrigger with query", + payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeContactPropertyTrigger, + ContactPropertyTrigger: &WorkflowContactPropertyTriggerPayload{ + ContactPropertyQuery: &WorkflowContactPropertyQuery{ + Key: "plan", + Is: WorkflowContactPropertyComparison{Value: WorkflowContactPropertyValue{String: ptr("pro")}, Operator: "equal"}, + }, + }, + }, + wantSubstr: []string{`"typeName":"ContactPropertyTrigger"`, `"key":"plan"`}, + }, + { + name: "AddToListTrigger", + payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeAddToListTrigger, + AddToListTrigger: &WorkflowAddToListTriggerPayload{ReEligible: ptr(false)}, + }, + wantSubstr: []string{`"typeName":"AddToListTrigger"`, `"reEligible":false`}, + }, + { + name: "AudienceFilter has no typeName", + payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeAudienceFilter, + AudienceFilter: &WorkflowAudienceFilterPayload{ + AudienceSegmentID: ptr("seg_1"), + AppliesDownstream: ptr(true), + }, + }, + wantSubstr: []string{`"audienceSegmentId":"seg_1"`, `"appliesDownstream":true`}, + notSubstr: []string{`typeName`}, + }, + { + name: "TimerAction has no typeName", + payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeTimerAction, + TimerAction: &WorkflowTimerActionPayload{Amount: ptr(3.0), Unit: WorkflowTimerUnitDays}, + }, + wantSubstr: []string{`"amount":3`, `"unit":"d"`}, + notSubstr: []string{`typeName`}, + }, + { + name: "ExperimentBranch has no typeName", + payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeExperimentBranchNode, + ExperimentBranch: &WorkflowExperimentBranchPayload{SamplingRate: ptr(50.0)}, + }, + wantSubstr: []string{`"samplingRate":50`}, + notSubstr: []string{`typeName`}, + }, + { + name: "Variant has no typeName", + payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeVariantNode, + Variant: &WorkflowVariantPayload{IsControl: ptr(true)}, + }, + wantSubstr: []string{`"isControl":true`}, + notSubstr: []string{`typeName`}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + raw, err := json.Marshal(tt.payload) + if err != nil { + t.Fatalf("Marshal: %v", err) + } + for _, want := range tt.wantSubstr { + if !strings.Contains(string(raw), want) { + t.Errorf("payload %s missing %q", raw, want) + } + } + for _, no := range tt.notSubstr { + if strings.Contains(string(raw), no) { + t.Errorf("payload %s must not contain %q", raw, no) + } + } + }) + } +} + +func TestUpdateWorkflowNode_RequestAndResponse(t *testing.T) { + var rawBody string + var gotBody map[string]any + resp := `{ + "id": "node_t", + "typeName": "TimerAction", + "nextNodeIds": ["n2"], + "amount": 2, + "unit": "h", + "workflowRevisionId": "rev_9" + }` + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + buf := new(strings.Builder) + tee := io.TeeReader(r.Body, buf) + json.NewDecoder(tee).Decode(&gotBody) + rawBody = buf.String() + if want := "/workflows/wf_1/nodes/node_t"; r.URL.Path != want { + t.Errorf("path = %q, want %q", r.URL.Path, want) + } + w.WriteHeader(http.StatusOK) + w.Write([]byte(resp)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + node, err := client.UpdateWorkflowNode("wf_1", "node_t", UpdateWorkflowNodeRequest{ + ExpectedRevisionID: ptr("rev_8"), + Payload: UpdateWorkflowNodePayload{ + TypeName: WorkflowNodeTypeTimerAction, + TimerAction: &WorkflowTimerActionPayload{Amount: ptr(2.0), Unit: WorkflowTimerUnitHours}, + }, + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if gotBody["expectedRevisionId"] != "rev_8" { + t.Errorf("expectedRevisionId = %v", gotBody["expectedRevisionId"]) + } + if !strings.Contains(rawBody, `"payload":{`) { + t.Errorf("payload missing from body: %s", rawBody) + } + if node.TypeName != WorkflowNodeTypeTimerAction || node.TimerAction == nil { + t.Errorf("node = %+v", node) + } + if node.TimerAction.Amount != 2 || node.TimerAction.Unit != WorkflowTimerUnitHours { + t.Errorf("timer = %+v", node.TimerAction) + } + if node.WorkflowRevisionID != "rev_9" { + t.Errorf("revision = %q, want rev_9", node.WorkflowRevisionID) + } +} + +func TestAddWorkflowBranch(t *testing.T) { + var gotBody map[string]any + resp := `{ + "node": {"id":"node_b","typeName":"BranchNode","nextNodeIds":["c1","c2","c3"],"workflowRevisionId":"rev_11"}, + "workflow": {"id":"wf_1","status":"Draft","workflowRevisionId":"rev_11","mailingListId":null,"rootNodeId":"r","nodes":{}} + }` + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotBody = decodeBody(t, r) + if want := "/workflows/wf_1/nodes/node_b/add-branch"; r.URL.Path != want { + t.Errorf("path = %q, want %q", r.URL.Path, want) + } + w.WriteHeader(http.StatusOK) + w.Write([]byte(resp)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + res, err := client.AddWorkflowBranch("wf_1", "node_b", AddWorkflowBranchRequest{ + ExpectedRevisionID: ptr("rev_10"), + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if gotBody["expectedRevisionId"] != "rev_10" { + t.Errorf("body = %+v", gotBody) + } + if res.Node.TypeName != WorkflowNodeTypeBranchNode || res.Node.BranchNode == nil { + t.Errorf("node = %+v", res.Node) + } + if len(res.Node.BranchNode.NextNodeIDs) != 3 { + t.Errorf("NextNodeIDs = %v", res.Node.BranchNode.NextNodeIDs) + } + if res.Node.WorkflowRevisionID != "rev_11" { + t.Errorf("revision = %q", res.Node.WorkflowRevisionID) + } + if res.Workflow.ID != "wf_1" { + t.Errorf("workflow = %+v", res.Workflow) + } +} + +func TestDeleteWorkflowNode(t *testing.T) { + var gotMethod, gotPath string + var gotBody map[string]any + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotMethod = r.Method + gotPath = r.URL.Path + gotBody = decodeBody(t, r) + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"status":"deleted","nodeIds":["node_x"],"workflowRevisionId":"rev_20","queuedContactCount":0}`)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + res, err := client.DeleteWorkflowNode("wf_1", "node_x", DeleteWorkflowNodeRequest{ + ExpectedRevisionID: ptr("rev_19"), + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if gotMethod != http.MethodDelete { + t.Errorf("method = %q, want DELETE", gotMethod) + } + if gotPath != "/workflows/wf_1/nodes/node_x" { + t.Errorf("path = %q", gotPath) + } + if gotBody["expectedRevisionId"] != "rev_19" { + t.Errorf("body = %+v", gotBody) + } + if res.Status != WorkflowMutationStatusDeleted { + t.Errorf("Status = %q, want deleted", res.Status) + } + if res.WorkflowRevisionID == nil || *res.WorkflowRevisionID != "rev_20" { + t.Errorf("WorkflowRevisionID = %v, want rev_20", res.WorkflowRevisionID) + } + if len(res.NodeIDs) != 1 || res.NodeIDs[0] != "node_x" { + t.Errorf("NodeIDs = %v", res.NodeIDs) + } +} + +func TestDeleteWorkflowNodeRecursive_DryRun(t *testing.T) { + var gotPath string + var gotBody map[string]any + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotPath = r.URL.Path + gotBody = decodeBody(t, r) + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"status":"dryRun","nodeIds":["node_x","node_y"],"queuedContactCount":4}`)) + })) + defer server.Close() + + client := NewClient("test-key", WithBaseURL(server.URL)) + res, err := client.DeleteWorkflowNodeRecursive("wf_1", "node_x", DeleteWorkflowNodeRequest{ + ExpectedRevisionID: ptr("rev_1"), + DryRun: true, + QueuedContactPolicy: WorkflowQueuedContactPolicyFail, + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if gotPath != "/workflows/wf_1/nodes/node_x/recursive" { + t.Errorf("path = %q", gotPath) + } + if gotBody["dryRun"] != true || gotBody["queuedContactPolicy"] != "fail" { + t.Errorf("body = %+v", gotBody) + } + if res.Status != WorkflowMutationStatusDryRun { + t.Errorf("Status = %q, want dryRun", res.Status) + } + if res.WorkflowRevisionID != nil { + t.Errorf("WorkflowRevisionID = %v, want nil", res.WorkflowRevisionID) + } + if len(res.NodeIDs) != 2 { + t.Errorf("NodeIDs = %v", res.NodeIDs) + } +} + +func TestWorkflowMutationNode_UnmarshalVariants(t *testing.T) { + tests := []struct { + name string + body string + want string + }{ + {"AddToListTrigger", `{"id":"n","typeName":"AddToListTrigger","nextNodeIds":[],"mailingListId":"ml_1","reEligible":true}`, WorkflowNodeTypeAddToListTrigger}, + {"AudienceFilter", `{"id":"n","typeName":"AudienceFilter","nextNodeIds":[],"appliesDownstream":true}`, WorkflowNodeTypeAudienceFilter}, + {"Variant", `{"id":"n","typeName":"VariantNode","nextNodeIds":[],"isControl":true}`, WorkflowNodeTypeVariantNode}, + {"ExperimentBranch", `{"id":"n","typeName":"ExperimentBranchNode","nextNodeIds":[],"samplingRate":25}`, WorkflowNodeTypeExperimentBranchNode}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + var n WorkflowMutationNode + if err := json.Unmarshal([]byte(tt.body), &n); err != nil { + t.Fatalf("Unmarshal: %v", err) + } + if n.TypeName != tt.want { + t.Errorf("TypeName = %q, want %q", n.TypeName, tt.want) + } + }) + } + + var af WorkflowMutationNode + if err := json.Unmarshal([]byte(`{"id":"n","typeName":"AddToListTrigger","nextNodeIds":[],"mailingListId":"ml_1","reEligible":true}`), &af); err != nil { + t.Fatalf("Unmarshal: %v", err) + } + if af.AddToListTrigger == nil || af.AddToListTrigger.MailingListID == nil || *af.AddToListTrigger.MailingListID != "ml_1" { + t.Errorf("AddToListTrigger = %+v", af.AddToListTrigger) + } +} + func TestWorkflowContactPropertyValue_RoundTrip(t *testing.T) { tests := []struct { name string @@ -516,3 +1141,92 @@ func TestWorkflowContactPropertyValue_RoundTrip(t *testing.T) { }) } } + +func TestUpdateWorkflowNodePayload_NilVariantErrors(t *testing.T) { + // A selected config variant with a nil pointer must error, not silently + // emit "null" — symmetric with the trigger variants. + for _, p := range []UpdateWorkflowNodePayload{ + {TypeName: WorkflowNodeTypeAudienceFilter}, + {TypeName: WorkflowNodeTypeTimerAction}, + {TypeName: WorkflowNodeTypeExperimentBranchNode}, + {TypeName: WorkflowNodeTypeVariantNode}, + } { + if _, err := json.Marshal(p); err == nil { + t.Errorf("%s: expected error marshaling nil variant, got none", p.TypeName) + } + } +} + +func TestCreateWorkflow_NullMailingList(t *testing.T) { + var rawBody string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + b, _ := io.ReadAll(r.Body) + rawBody = string(b) + w.Write([]byte(`{"id":"wf_1","status":"Draft","workflowRevisionId":"rev_1","mailingListId":null,"rootNodeId":"r","nodes":{}}`)) + })) + defer server.Close() + + client := NewClient("key", WithBaseURL(server.URL)) + if _, err := client.CreateWorkflow(CreateWorkflowRequest{Name: "New flow"}); err != nil { + t.Fatalf("CreateWorkflow: %v", err) + } + if !strings.Contains(rawBody, `"mailingListId":null`) { + t.Errorf("body must send explicit null mailingListId when unset: %s", rawBody) + } +} + +func TestWorkflowNodeWithRevision_MarshalKeepsRevision(t *testing.T) { + n := WorkflowNodeWithRevision{ + WorkflowNode: WorkflowNode{ + TypeName: WorkflowNodeTypeTimerAction, + TimerAction: &TimerActionWorkflowNode{ID: "n1", WorkflowID: "wf_1", NextNodeIDs: []string{}, Amount: 5, Unit: WorkflowTimerUnitHours}, + }, + WorkflowRevisionID: ptr("rev_9"), + } + raw, err := json.Marshal(n) + if err != nil { + t.Fatalf("Marshal: %v", err) + } + if !strings.Contains(string(raw), `"workflowRevisionId":"rev_9"`) { + t.Errorf("marshal dropped workflowRevisionId: %s", raw) + } + if !strings.Contains(string(raw), `"typeName":"TimerAction"`) { + t.Errorf("marshal dropped node fields: %s", raw) + } +} + +func TestWorkflowMutationNodeWithRevision_MarshalKeepsRevision(t *testing.T) { + const in = `{"id":"n","typeName":"AddToListTrigger","nextNodeIds":[],"mailingListId":"ml_1","reEligible":true,"workflowRevisionId":"rev_7"}` + var n WorkflowMutationNodeWithRevision + if err := json.Unmarshal([]byte(in), &n); err != nil { + t.Fatalf("Unmarshal: %v", err) + } + raw, err := json.Marshal(n) + if err != nil { + t.Fatalf("Marshal: %v", err) + } + if !strings.Contains(string(raw), `"workflowRevisionId":"rev_7"`) { + t.Errorf("marshal dropped workflowRevisionId: %s", raw) + } + if !strings.Contains(string(raw), `"typeName":"AddToListTrigger"`) { + t.Errorf("marshal dropped node fields: %s", raw) + } +} + +func TestCreatedWorkflowNode_MarshalKeepsRevisionAndChildren(t *testing.T) { + const in = `{"id":"n","typeName":"AddToListTrigger","nextNodeIds":[],"mailingListId":"ml_1","reEligible":true,"workflowRevisionId":"rev_7","createdChildNodes":[{"id":"c","typeName":"AddToListTrigger","nextNodeIds":[],"mailingListId":"ml_2","reEligible":false}]}` + var n CreatedWorkflowNode + if err := json.Unmarshal([]byte(in), &n); err != nil { + t.Fatalf("Unmarshal: %v", err) + } + raw, err := json.Marshal(n) + if err != nil { + t.Fatalf("Marshal: %v", err) + } + if !strings.Contains(string(raw), `"workflowRevisionId":"rev_7"`) { + t.Errorf("marshal dropped workflowRevisionId: %s", raw) + } + if !strings.Contains(string(raw), `"createdChildNodes"`) { + t.Errorf("marshal dropped createdChildNodes: %s", raw) + } +}