diff --git a/framework/clclient/client.go b/framework/clclient/client.go index 1f1b63258..90d1aa4bc 100644 --- a/framework/clclient/client.go +++ b/framework/clclient/client.go @@ -9,6 +9,7 @@ import ( "net/http" "os" "regexp" + "strconv" "strings" "sync" "time" @@ -1431,3 +1432,38 @@ func ImportP2PKeys(cl []*ChainlinkClient, keys [][]byte) error { } return eg.Wait() } + +func (c *ChainlinkClient) ReadWorkflowEvents(workflowID string, sequence int64, limit int) (*WorkflowDebugEvents, *http.Response, error) { + specObj := &WorkflowDebugEvents{} + framework.L.Info().Str(NodeURL, c.Config.URL).Str("ID", workflowID).Int64("sequence", sequence).Int("limit", limit).Msg("Reading Workflow Events") + resp, err := c.APIClient.R(). + SetResult(&specObj). + SetPathParams(map[string]string{ + "id": workflowID, + }). + SetQueryParams(map[string]string{ + "sequence": strconv.FormatInt(sequence, 10), + "limit": strconv.Itoa(limit), + }). + Get("/v2/debug/workflow/{id}/events?sequence={sequence}&limit={limit}") + if err != nil { + return nil, nil, err + } + return specObj, resp.RawResponse, err +} + +func (c *ChainlinkClient) ReadOrphanEvents(sequence int64, limit int) (*WorkflowOrphanEvents, *http.Response, error) { + specObj := &WorkflowOrphanEvents{} + framework.L.Info().Str(NodeURL, c.Config.URL).Int64("sequence", sequence).Int("limit", limit).Msg("Reading Workflow Orphan Events") + resp, err := c.APIClient.R(). + SetResult(&specObj). + SetQueryParams(map[string]string{ + "sequence": strconv.FormatInt(sequence, 10), + "limit": strconv.Itoa(limit), + }). + Get("/v2/debug/workflow/orphan_events?sequence={sequence}&limit={limit}") + if err != nil { + return nil, nil, err + } + return specObj, resp.RawResponse, err +} diff --git a/framework/clclient/models.go b/framework/clclient/models.go index 69f9d9d93..5d8eaa794 100644 --- a/framework/clclient/models.go +++ b/framework/clclient/models.go @@ -1419,3 +1419,45 @@ type ForwarderAttributes struct { CreatedAt time.Time `json:"createdAt"` UpdatedAt time.Time `json:"updatedAt"` } + +type WorkflowDebugEvents struct { + Data WorkflowDebugEventsData `json:"data"` +} + +type WorkflowDebugEventsData struct { + Type string `json:"type"` + ID string `json:"id"` + Attributes WorkflowDebugEventsAttributes `json:"attributes"` +} + +type WorkflowDebugEventsAttributes struct { + Events []WorkflowDebugEvent `json:"events"` +} + +type WorkflowDebugEvent struct { + Timestamp time.Time `json:"timestamp"` + Sequence int64 `json:"sequence"` + Message []byte `json:"message"` // protobuf encoded + Type string `json:"type"` // protobuf type name +} + +type WorkflowOrphanEvents struct { + Data WorkflowOrphanEventsData `json:"data"` +} + +type WorkflowOrphanEventsData struct { + Type string `json:"type"` + ID string `json:"id"` + Attributes WorkflowDebugEventsAttributes `json:"attributes"` +} + +type WorkflowOrphanEventsAttributes struct { + Events []WorkflowOrphanEvent `json:"events"` +} + +type WorkflowOrphanEvent struct { + Timestamp time.Time `json:"timestamp"` + Sequence int64 `json:"sequence"` + Message []byte `json:"message"` // protobuf encoded + Type string `json:"type"` // protobuf type name +}