Skip to content

Commit de29855

Browse files
committed
feat(subscriptions): Scripts for consuming events
1 parent b9f9c18 commit de29855

5 files changed

Lines changed: 574 additions & 11 deletions

File tree

‎test/utils/dump.go‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,13 @@ import (
99

1010
// DumpJSON dumps the given value as JSON to the given writer.
1111
func DumpJSON(v any, w io.Writer) {
12+
if bytesData, ok := v.([]byte); ok {
13+
jsonData := make(map[string]any)
14+
if err := json.Unmarshal(bytesData, &jsonData); err == nil {
15+
v = jsonData
16+
}
17+
}
18+
1219
// Convert any error interfaces recursively before encoding.
1320
convertedValue := substituteErrorsToStrings(v)
1421

‎test/utils/testscenario/crud.go‎

Lines changed: 34 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import (
44
"context"
55
"errors"
66
"fmt"
7+
"os"
78
"strconv"
89
"time"
910

@@ -61,6 +62,10 @@ func (p Property) isZero() bool {
6162
return p.Key == "" && p.Value == ""
6263
}
6364

65+
func (p Property) String() string {
66+
return fmt.Sprintf("{ %v : %v }", p.Key, p.Value)
67+
}
68+
6469
// ValidateCreateUpdateDelete is a comprehensive test scenario utilizing Read/Write/Delete connector operations.
6570
//
6671
// Flow:
@@ -69,7 +74,8 @@ func (p Property) isZero() bool {
6974
// 3. Update the object using the "UP" payload.
7075
// 4. Read again and verify updates took effect.
7176
// 5. Delete the object at the end.
72-
func ValidateCreateUpdateDelete[CP, UP any](ctx context.Context, conn ConnectorCRUD, objectName string,
77+
func ValidateCreateUpdateDelete[CP, UP any](
78+
ctx context.Context, conn ConnectorCRUD, objectName string,
7379
createPayload CP, updatePayload UP, suite CRUDTestSuite,
7480
) {
7581
fmt.Println("> TEST Create/Update/Delete", objectName)
@@ -114,7 +120,7 @@ func ValidateCreateUpdateDelete[CP, UP any](ctx context.Context, conn ConnectorC
114120

115121
// UPDATE
116122
fmt.Println("Updating some object properties")
117-
err = updateObject(ctx, conn, objectName, objectID, &updatePayload)
123+
_, err = updateObject(ctx, conn, objectName, objectID, &updatePayload)
118124
failOnError(err)
119125
fmt.Println("Validate object has changed accordingly")
120126

@@ -181,7 +187,7 @@ func getRecordIdentifierValue(object map[string]any, key string) string {
181187
}
182188

183189
func createObject[CP any](
184-
ctx context.Context, conn ConnectorCRUD, objectName string, payload *CP,
190+
ctx context.Context, conn connectors.WriteConnector, objectName string, payload *CP,
185191
) (*common.WriteResult, error) {
186192
res, err := conn.Write(ctx, common.WriteParams{
187193
ObjectName: objectName,
@@ -199,8 +205,10 @@ func createObject[CP any](
199205
return res, nil
200206
}
201207

202-
func readObjects(ctx context.Context, conn ConnectorCRUD,
203-
objectName string, fields datautils.StringSet, since time.Time) (*common.ReadResult, error) {
208+
func readObjects(
209+
ctx context.Context, conn connectors.ReadConnector,
210+
objectName string, fields datautils.StringSet, since time.Time,
211+
) (*common.ReadResult, error) {
204212
res, err := conn.Read(ctx, common.ReadParams{
205213
ObjectName: objectName,
206214
Fields: fields,
@@ -247,22 +255,22 @@ func searchObjectRecord(res *common.ReadResult, key, value string) (*objectRecor
247255
}
248256

249257
func updateObject[UP any](
250-
ctx context.Context, conn ConnectorCRUD, objectName string, objectID string, payload *UP,
251-
) error {
258+
ctx context.Context, conn connectors.WriteConnector, objectName string, objectID string, payload *UP,
259+
) (*common.WriteResult, error) {
252260
res, err := conn.Write(ctx, common.WriteParams{
253261
ObjectName: objectName,
254262
RecordId: objectID,
255263
RecordData: payload,
256264
})
257265
if err != nil {
258-
return fmt.Errorf("error updating object: %w", err)
266+
return nil, fmt.Errorf("error updating object: %w", err)
259267
}
260268

261269
if !res.Success {
262-
return errors.New("failed to update an object")
270+
return nil, errors.New("failed to update an object")
263271
}
264272

265-
return nil
273+
return res, nil
266274
}
267275

268276
func removeObject(ctx context.Context, conn ConnectorCRUD, objectName string, objectID string) error {
@@ -283,6 +291,21 @@ func removeObject(ctx context.Context, conn ConnectorCRUD, objectName string, ob
283291

284292
func failOnError(err error) {
285293
if err != nil {
286-
utils.Fail("fatal", "error", err)
294+
utils.Fail("[test failed]", "error", err)
287295
}
288296
}
297+
298+
// printError prints error and returns true if error is not nil.
299+
func printError(err error) bool {
300+
if err == nil {
301+
return false
302+
}
303+
304+
fmt.Println("[test failed]", "error", err.Error())
305+
if httpError, ok := errors.AsType[*common.HTTPError](err); ok {
306+
utils.DumpJSON(httpError.Body, os.Stdout)
307+
return true
308+
}
309+
310+
return true
311+
}
Lines changed: 217 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,217 @@
1+
package testscenario
2+
3+
import (
4+
"context"
5+
"fmt"
6+
"os"
7+
"time"
8+
9+
"github.com/amp-labs/connectors"
10+
"github.com/amp-labs/connectors/common"
11+
"github.com/amp-labs/connectors/internal/components"
12+
"github.com/amp-labs/connectors/internal/datautils"
13+
"github.com/amp-labs/connectors/test/utils"
14+
)
15+
16+
type ConnectorWebhookSubscriber interface {
17+
components.SubscriptionCreator
18+
components.WebhookMessageVerifier
19+
connectors.ReadConnector
20+
connectors.WriteConnector
21+
connectors.DeleteConnector
22+
}
23+
24+
type SubscribeParamBuilder func(webhookURL string) *common.SubscribeParams
25+
26+
type SubscribeReceiveEventsSuite struct {
27+
SubscribeParamBuilder SubscribeParamBuilder
28+
// ExpectedWebhookCalls is the number of events to wait before quiting the script.
29+
//
30+
// Number of Operations may be greater than the number of expected events.
31+
// Some operations can be utilized for cleaning records.
32+
// Ex: You are testing only update. The operations would still need create, delete, not just update.
33+
//
34+
// Webhook may be called more if there are already some subscriptions that would intervene.
35+
// There is no way to protect against these side effects.
36+
ExpectedWebhookCalls int
37+
Operations []ConnectorOperation
38+
// WebhookRouter allows a list of alternative request handling by the webhook handler.
39+
// When conditions are met for the Route that handler is executed.
40+
// If there are no custom routes or none match then default webhook handling will take place.
41+
// This includes printing events until reaching the number of ExpectedWebhookCalls or
42+
// script cancellation.
43+
WebhookRouter WebhookRouter
44+
VerificationParams *common.VerificationParams
45+
}
46+
47+
type ConnectorMethod string
48+
49+
const (
50+
ConnectorMethodCreate ConnectorMethod = "create"
51+
ConnectorMethodUpdate ConnectorMethod = "update"
52+
ConnectorMethodDelete ConnectorMethod = "delete"
53+
)
54+
55+
type ConnectorOperation struct {
56+
// ObjectName object to create, update or to remove.
57+
ObjectName string
58+
// Method invokes Write() or Delete().
59+
Method ConnectorMethod
60+
// Payload relevant for ConnectorMethodCreate and ConnectorMethodUpdate.
61+
Payload any
62+
// SearchProcedure relevant for ConnectorMethodUpdate and ConnectorMethodDelete.
63+
SearchProcedure SearchProcedure
64+
}
65+
66+
type SearchProcedure struct {
67+
ReadFields datautils.StringSet
68+
RecordIdentifierKey string
69+
WaitBeforeSearch time.Duration
70+
SearchBy Property
71+
}
72+
73+
// ValidateSubscribeReceiveEvents is a comprehensive test scenario utilizing subscription connector operations.
74+
//
75+
// Flow:
76+
// 1. Starts local server
77+
// 2. Asks user for public URL (ngrok)
78+
// 3. Creates subscription
79+
// 4. Optionally triggers events (Write)
80+
// 5. Waits for webhook(s)
81+
// 6. Exits cleanly
82+
func ValidateSubscribeReceiveEvents(
83+
ctx context.Context,
84+
conn ConnectorWebhookSubscriber,
85+
suite SubscribeReceiveEventsSuite,
86+
) {
87+
fmt.Println("> TEST Subscribe/Write/Recieve")
88+
89+
fmt.Println("============== Starting Webhook Handler ==================")
90+
messageChannel := make(chan webhookMessageResult)
91+
webhookURL, shutdown := startWebhookHandler(ctx, conn,
92+
suite.WebhookRouter, suite.VerificationParams, messageChannel,
93+
)
94+
defer shutdown()
95+
96+
fmt.Printf("Local webhook server started at: \"%s\"\n", webhookURL)
97+
publicURL, ok := getPublicWebhookURL(ctx)
98+
if !ok {
99+
return
100+
}
101+
102+
fmt.Println("============== Invoking connector.Subscribe() ==================")
103+
params := *suite.SubscribeParamBuilder(publicURL)
104+
subscriptionResult, err := conn.Subscribe(ctx, params)
105+
if printError(err) {
106+
return
107+
}
108+
109+
switch subscriptionResult.Status {
110+
case common.SubscriptionStatusPending:
111+
fmt.Printf("Connector returned status (%v). Script is not designed to handle this.\n",
112+
common.SubscriptionStatusPending)
113+
return
114+
case common.SubscriptionStatusFailed:
115+
utils.DumpJSON(subscriptionResult.Result, os.Stdout)
116+
fmt.Println("Subscription failed ❌")
117+
return
118+
case common.SubscriptionStatusSuccess:
119+
utils.DumpJSON(subscriptionResult.Result, os.Stdout)
120+
utils.DumpJSON(subscriptionResult.ObjectEvents, os.Stdout)
121+
// continue script execution.
122+
case common.SubscriptionStatusFailedToRollback:
123+
fmt.Println("Subscription encountered failures and then failed to rollback ❌")
124+
utils.DumpJSON(subscriptionResult.Result, os.Stdout)
125+
utils.DumpJSON(subscriptionResult.ObjectEvents, os.Stdout)
126+
return
127+
default:
128+
fmt.Printf("Unknown subscription status (%v)\n", subscriptionResult.Status)
129+
return
130+
}
131+
132+
fmt.Println("============== Invoking connector.Write/Delete() ==================")
133+
for _, trigger := range suite.Operations {
134+
switch trigger.Method {
135+
case ConnectorMethodCreate:
136+
fmt.Printf("Creating object %v:\n", trigger.ObjectName)
137+
createResult, err := createObject[any](ctx, conn, trigger.ObjectName, &trigger.Payload)
138+
if printError(err) {
139+
return
140+
}
141+
utils.DumpJSON(createResult, os.Stdout)
142+
case ConnectorMethodUpdate:
143+
objectID, ok := searchForRecord(ctx, conn, trigger.ObjectName, trigger.SearchProcedure)
144+
if !ok {
145+
return
146+
}
147+
148+
fmt.Printf("Updating object %v:\n", trigger.ObjectName)
149+
updateResult, err := updateObject[any](ctx, conn, trigger.ObjectName, objectID, &trigger.Payload)
150+
if printError(err) {
151+
return
152+
}
153+
utils.DumpJSON(updateResult, os.Stdout)
154+
case ConnectorMethodDelete:
155+
objectID, ok := searchForRecord(ctx, conn, trigger.ObjectName, trigger.SearchProcedure)
156+
if !ok {
157+
return
158+
}
159+
160+
fmt.Printf("Deleting object %v:\n", trigger.ObjectName)
161+
err = removeObject(ctx, conn, trigger.ObjectName, objectID)
162+
if printError(err) {
163+
return
164+
}
165+
fmt.Println("... object deleted.")
166+
}
167+
}
168+
169+
// Waiting for the events to arrive. Then report on them and exit.
170+
// This can be stopped prematurely via context cancellation.
171+
receivedNumEvents := 0
172+
fmt.Printf("============== Waiting for %d webhook messages ==================\n", suite.ExpectedWebhookCalls)
173+
174+
for receivedNumEvents < suite.ExpectedWebhookCalls {
175+
select {
176+
case message := <-messageChannel:
177+
receivedNumEvents++
178+
fmt.Printf("[%d/%d] Received webhook message:\n", receivedNumEvents, suite.ExpectedWebhookCalls)
179+
if message.Error == "" {
180+
utils.DumpJSON(message.Body, os.Stdout)
181+
} else {
182+
utils.DumpJSON(message.Error, os.Stdout)
183+
}
184+
185+
case <-ctx.Done():
186+
fmt.Println("Context cancelled, stopping...")
187+
return
188+
}
189+
}
190+
191+
fmt.Println("============== Done ==================")
192+
}
193+
194+
func searchForRecord(
195+
ctx context.Context, conn ConnectorWebhookSubscriber, objectName string, procedure SearchProcedure,
196+
) (string, bool) {
197+
if procedure.WaitBeforeSearch != 0 {
198+
fmt.Println("... waiting")
199+
time.Sleep(procedure.WaitBeforeSearch)
200+
}
201+
202+
fmt.Printf("Search object %v by %v\n", objectName, procedure.SearchBy.String())
203+
res, err := readObjects(ctx, conn, objectName, procedure.ReadFields, procedure.SearchBy.Since)
204+
if printError(err) {
205+
return "", false
206+
}
207+
208+
search := procedure.SearchBy
209+
object, err := searchObjectRecord(res, search.Key, search.Value)
210+
if printError(err) {
211+
return "", false
212+
}
213+
214+
objectID := object.getRecordIdentifierValue(procedure.RecordIdentifierKey)
215+
216+
return objectID, true
217+
}

0 commit comments

Comments
 (0)