/
retrieveResource.go
82 lines (76 loc) · 3.35 KB
/
retrieveResource.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
package service
import (
"bytes"
"context"
"fmt"
"net/http"
"time"
"github.com/plgd-dev/hub/v2/cloud2cloud-connector/events"
"github.com/plgd-dev/hub/v2/cloud2cloud-connector/store"
"github.com/plgd-dev/hub/v2/pkg/log"
kitNetGrpc "github.com/plgd-dev/hub/v2/pkg/net/grpc"
"github.com/plgd-dev/hub/v2/resource-aggregate/commands"
raEvents "github.com/plgd-dev/hub/v2/resource-aggregate/events"
raService "github.com/plgd-dev/hub/v2/resource-aggregate/service"
"go.opentelemetry.io/otel/trace"
)
func retrieveDeviceResource(ctx context.Context, traceProvider trace.TracerProvider, deviceID, href string, linkedAccount store.LinkedAccount, linkedCloud store.LinkedCloud) (string, []byte, commands.Status, error) {
client := linkedCloud.GetHTTPClient(traceProvider)
defer client.CloseIdleConnections()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, makeHTTPEndpoint(linkedCloud.Endpoint.URL, deviceID, href), nil)
if err != nil {
return "", nil, commands.Status_BAD_REQUEST, fmt.Errorf("cannot create post request: %w", err)
}
req.Header.Set(AcceptHeader, events.ContentType_JSON+","+events.ContentType_VNDOCFCBOR)
req.Header.Set(AuthorizationHeader, AuthorizationBearerPrefix+string(linkedAccount.Data.Target().AccessToken))
req.Header.Set("Connection", "close")
req.Close = true
httpResp, err := client.Do(req)
if err != nil {
return "", nil, commands.Status_UNAVAILABLE, fmt.Errorf("cannot post: %w", err)
}
defer func() {
if errC := httpResp.Body.Close(); errC != nil {
log.Errorf("failed to close response body stream: %w", errC)
}
}()
if httpResp.StatusCode != http.StatusOK {
status := commands.HTTPStatus2Status(httpResp.StatusCode)
return "", nil, status, fmt.Errorf("unexpected statusCode %v", httpResp.StatusCode)
}
respContentType := httpResp.Header.Get(events.ContentTypeKey)
respContent := bytes.NewBuffer(make([]byte, 0, 1024))
_, err = respContent.ReadFrom(httpResp.Body)
if err != nil {
return "", nil, commands.Status_UNAVAILABLE, fmt.Errorf("cannot read update response: %w", err)
}
return respContentType, respContent.Bytes(), commands.Status_OK, nil
}
func retrieveResource(ctx context.Context, traceProvider trace.TracerProvider, raClient raService.ResourceAggregateClient, e *raEvents.ResourceRetrievePending, linkedAccount store.LinkedAccount, linkedCloud store.LinkedCloud) error {
deviceID := e.GetResourceId().GetDeviceId()
href := e.GetResourceId().GetHref()
contentType, content, status, err := retrieveDeviceResource(ctx, traceProvider, deviceID, href, linkedAccount, linkedCloud)
if err != nil {
log.Errorf("cannot update resource /%v%v: %w", deviceID, href, err)
}
coapContentFormat := stringToSupportedMediaType(contentType)
ctx = kitNetGrpc.CtxWithToken(ctx, linkedAccount.Data.Origin().AccessToken.String())
_, err = raClient.ConfirmResourceRetrieve(ctx, &commands.ConfirmResourceRetrieveRequest{
ResourceId: commands.NewResourceID(deviceID, href),
CorrelationId: e.GetAuditContext().GetCorrelationId(),
CommandMetadata: &commands.CommandMetadata{
ConnectionId: linkedAccount.ID,
Sequence: uint64(time.Now().UnixNano()),
},
Content: &commands.Content{
Data: content,
ContentType: contentType,
CoapContentFormat: coapContentFormat,
},
Status: status,
})
if err != nil {
return fmt.Errorf("cannot update resource /%v%v: %w", deviceID, href, err)
}
return nil
}