diff --git a/api/queryService/api.gen.go b/api/queryService/api.gen.go index 1a0ad992..2a9b0846 100644 --- a/api/queryService/api.gen.go +++ b/api/queryService/api.gen.go @@ -65,6 +65,18 @@ type ClientUpdate struct { Name *string `json:"name,omitempty"` } +// Document defines model for Document. +type Document struct { + // Bucket The bucket containing the document + Bucket string `json:"bucket"` + + // Id The document id + Id openapi_types.UUID `json:"id"` + + // Key The path to the document + Key string `json:"key"` +} + // ExportDetails Payload for export trigger response. type ExportDetails struct { // JobId The job id relative to the export. @@ -206,6 +218,9 @@ type JobUpdate struct { CanSync *bool `json:"can_sync,omitempty"` } +// ListDocuments The documents in the job. +type ListDocuments = []Document + // ListQueries defines model for ListQueries. type ListQueries struct { // Queries List of queries. @@ -331,6 +346,9 @@ type ServerInterface interface { // Update a job collector // (PATCH /job/{id}/collector) UpdateJobCollectorByJobId(ctx echo.Context, id openapi_types.UUID) error + // List the documents for a job + // (GET /job/{id}/documents) + ListDocumentsByJobId(ctx echo.Context, id openapi_types.UUID) error // Check export state. // (GET /job/{id}/export) ExportState(ctx echo.Context, id openapi_types.UUID) error @@ -479,6 +497,22 @@ func (w *ServerInterfaceWrapper) UpdateJobCollectorByJobId(ctx echo.Context) err return err } +// ListDocumentsByJobId converts echo context to params. +func (w *ServerInterfaceWrapper) ListDocumentsByJobId(ctx echo.Context) error { + var err error + // ------------- Path parameter "id" ------------- + var id openapi_types.UUID + + err = runtime.BindStyledParameterWithOptions("simple", "id", ctx.Param("id"), &id, runtime.BindStyledParameterOptions{ParamLocation: runtime.ParamLocationPath, Explode: false, Required: true}) + if err != nil { + return echo.NewHTTPError(http.StatusBadRequest, fmt.Sprintf("Invalid format for parameter id: %s", err)) + } + + // Invoke the callback with all the unmarshaled arguments + err = w.Handler.ListDocumentsByJobId(ctx, id) + return err +} + // ExportState converts echo context to params. func (w *ServerInterfaceWrapper) ExportState(ctx echo.Context) error { var err error @@ -598,6 +632,7 @@ func RegisterHandlersWithBaseURL(router EchoRouter, si ServerInterface, baseURL router.PATCH(baseURL+"/job/:id", wrapper.UpdateJob) router.GET(baseURL+"/job/:id/collector", wrapper.GetJobCollectorByJobId) router.PATCH(baseURL+"/job/:id/collector", wrapper.UpdateJobCollectorByJobId) + router.GET(baseURL+"/job/:id/documents", wrapper.ListDocumentsByJobId) router.GET(baseURL+"/job/:id/export", wrapper.ExportState) router.GET(baseURL+"/query", wrapper.ListQueries) router.POST(baseURL+"/query", wrapper.CreateQuery) @@ -610,48 +645,51 @@ func RegisterHandlersWithBaseURL(router EchoRouter, si ServerInterface, baseURL // Base64 encoded, gzipped, json marshaled Swagger object var swaggerSpec = []string{ - "H4sIAAAAAAAC/+RbW3PbNhb+KxjuPmokt+mT3lLb6ciTjbOxstOdTkYDgkcSXJCgAVCO6tF/38GFNxEQ", - "KdtsvdOnxCZwDvCd71yAAz9FhKc5zyBTMpo/RZJsIcXmv5eMQqYuBWAF+udc8ByEomC+Zjg1v01AEkFz", - "RXkWzaPlFhAx85AZMInUPodoHkklaLaJDodJJOChoAKSaP6blfKtGsXjeyAqOkyc8q954lVOcLaS+4x0", - "F7BYI1WvgUqEGeOPNNsgTBTdAdLTZL2umHMGONMqn7+jzuqvv+dcqCtQmDLZlfkZ7xnHCVpzgcAMRUrQ", - "zQYEEiBznkmYRpOjPd/zeEUT/wLveYxoggQwbDapuEHBytai1lykWEXzqCho0t3DJOKFygu1YpxgK9en", - "pvyKaIYet5RsG1rQHzRHa8oAPVLGUAxozYss0crhO05zZvS9m89mT3FBfgd1mD1ZXGlymD3d89j8q2gK", - "UuE0P6yerGCaHKZ/0Ny3aKmwKgw4/xSwjubRP2Y1m2eOyjNrjDs79piADtVK1regNe8qZV1gci4ljVmF", - "hbaHFgjS7D8rUq1Lr42BAq2OZqtc8I0Aqcm4xpRB0lBe79EqX1p6nKaS45Ahe9YwfptHawosWa0pUyA8", - "2/lgPhjDSsJzQDGWkCCeITMRWaKgHWaF3R1VkPba4IOea0VHtcNgIfBe/0yzDUi9gOesq5qMcixwCnp+", - "d9uQJasymHiYjaVC+jPi61rgNEA6oYKiqERrKgYK84WOAW7ObYx7KEDsdcAomJLa5eOSf5AM8Hm/I/j4", - "3zRewC6Gc4hwVqRZB3rCs4SWUWUgSy6rOYeJo2w4PusvJSqWpTTT8ISNaOnrl2a/oULCumBlKLW8bPG9", - "I7TN6iN8G3uYNBCpltID/GUTwxMByG5/XRml0tSKRAykXKkt1vo3JsGL8kfCuIRkRTMFYodZNIl4Dlnz", - "ZwZrteoOE3Sz9f2eZoQVCZg8YP/nC3OL5F8gJd54kr3PG75m9KEARBOdP9YUhM2lmaJqfz73A7y/4fE5", - "pcddDkSvRaLHLaitC1bGZ6WrPtje1B96Db4KxObDoPuXRU1SCY6B8WyjnX9Ijj8dVoZIGJZwb3gcyLZG", - "ar3LSY3myQR8w2NbD/5JlWCvAYZg9ToFspFdRo1yryGMOGNAFPdE6eZXlPIEWDdDWlhWOxAyGGccdG5M", - "GXNJKbrlejRT736sN6hDwsamfxOkAtHXftMBvCR5S/qgaqO5WxNDfTXHqUS7uCq3hqXkhGIFiV7LoFKa", - "6bpPncbRjnkpjinNaFqkK6KJC+K0Sje40mkqxi2ghJMi1Xx0Us7UreC7ep5iPVOXbHKgzlDl7tgUhiOw", - "2I6lJscu0OdnllthGqPawapddxg99FjtCm8rWA/zFjamKAzy2paMJoFgZQ9qOc8LjYNZnC0mz0+gLkBV", - "yvuAq8/14dNMocfYurI59W8ft/52Tu9lUuBWatQK6oj0ta4A3UM3BnWhaE6zZpSmJC5TTFmmLz6t7v77", - "6TKaRJ9ul+a/11fe6vmGx+dflb20XvWZ5iOV6t8FCKe8vZSH+kN7JXqWBsANGOwvWtO+9+RVqvVZyUro", - "rPT5IcWEwIHeRHi2ppuu9Evz+0LYu7bScyrJg+p6/wHJv8LXL2TOQaE01aqXHuVIl8YWV22m9O7nOJLa", - "nwdQbKkHeitzI6JTN3SgC1IvGMfOokYGjyfo8f8KsJkdRG4JUn0Bc8nTha9Mb95EUBf3VRZUHBka4w2m", - "mRx2Y26LnaBz/MfjEVpPIcHaDaSubp5T+TZ3d7yOHrxsc6ELmLmE6m7ii6kHj/YAFqDTydEKDK/GMeNU", - "YtQzOwGlTI03d7efVte/Lr+8v1zefokm0eXtp+X1r8vVh68fP3pzpNEbypJ9Ad/OS46C/gvjfCmUnBnv", - "+x26lMxGcexu4j+YO/w1767k/edF5Wbt/Rl7oFtBtiCV27sEsaPEVgCKKtMz8o17/3kRTaLKWtEP04vp", - "hWlj5ZDhnEbz6N30YvpOnxWw2ppNzkh9hcRt0DiKrSYWS4RNPHWl4iNVtsuVC76jCSQosV29qb0ZtQta", - "JNV8d1FlrQRS/cyTvYvoyqnHec6o7aTN7qWlmw2HfcGy1ZI9tH1OiQLML6yPmz3/ePHDq+mub2iN4iPs", - "LFrErCxBsiAEpFwXjO2n2i4/XVx4AnG2w8y0LQ1SKOaJHn2YRLJIU6wLMwdqyyaaG3gjdRSwau8sa6Jv", - "eqoz8+yJJgetcgPKF9eUoLAzxpY26JDS4vEeUSXR4qpr4V9ANczbwvni1XCuLzvDOFccHAptszNm5vzk", - "4b+7luSqbN22TfELKIQbKC2uTlhCO16pMpr/dvqSzclUHAlnGBOe9EDtvuUN6NzWXG3KTxqw9p3fvplw", - "QLahgClt35SazNyKAJp9Qce3k/8Ex3fZa5DjXwQNXLjk8CInHcihIwrZDVQs6vHke9d/GRCt9an17FB9", - "w+ORzFXfULylIH3D47Ej9L2BtDTqDY87Fp3ZFnXYsIuMKmpM23hckguul9u1pHsYYV9JjGTN9hOMt2TR", - "69bjoZF8WtMmmBQcKvWDk4b53buZDgPOT83mpu5kXi59ebykHPKoUXLxScxtInaYtLJwy+EGpGAtwwL6", - "VyfbKoAPyLSjxu1zcmzAcGMm2JPEqNLrgDCsnXBGmh3j89yxmtrvmFVP4+f9DY8XydjVc7WpUA6slt72", - "3QDc9fABHtmC5dg3q5U9w0mrc3O7rfQG3Lbe8hkO7OXEOIXYUdPxJZ5db3VsHx9Auqa31+N7GdcKAHUx", - "5vX+yy2Q320hRgohdGlvnrWablXziWnbyvWDWRjT19uvrMPl0RZLFANkqHqAG8b+un7AGwTeoFJWpgaP", - "6Ymq56yTcI3oCJ6tLf9QNrx6oz1rN+Wa74DNS1wiqAJBcdf8zSbgiOZvqvEY//3xDl63QrPXkkGOmJbK", - "QwVCSQ4zq8WNAYdbe2l79vHW9jbHiavN5tVbOhBZq4x8yH1wwAaMWrnZMw461tSnK6rariM5lmurh9Ad", - "5cDT4062wKrQaRVWXZ8aHm+rJstfXUg1fHxAATW6a7/wptFac8waqYcvVW10lq/OVNlNfjMM8qaH6+9A", - "CpsfTNtaFJl9PxTKFO0/jDm6QQOpRudTs1s/nFSvrd91vz2RbWlgdK8fX5IojKB+2uk5IHZ+gn0WPClI", - "1RAFXdMXgkXzaKtULuezGc7p1P1V3ZTwdLb7IdJscdqO5d2W9pb2LwQhQYq7K3hZs7V9BX+YDBNzz+OG", - "jMZVwxkC6tNLW1Tn9DJUZl14OWEt/IdKscV4Q0q7tD98O/wvAAD//8+zA+A+OwAA", + "H4sIAAAAAAAC/+RbX2/bOBL/KoTuHg07u90nv3WTdOGg1/Sa9LCHRWFQ0thmKpMKSSX1Bv7uB/6VZJGW", + "nETdHvapTUTOcGZ+84czzFOSsW3JKFApkvlTIrINbLH+73lBgMpzDliC+rnkrAQuCeivFG/1b3MQGSel", + "JIwm8+R2AyjT+5BeMEnkroRkngjJCV0n+/0k4XBfEQ55Mv/DUPniV7H0DjKZ7CeW+ecyDzLPMF2KHc26", + "B1iskKzPQATCRcEeCV0jnEnyAEhtE/W5UsYKwFSxfL5EndNfsKzaApXdk6dV9hVkmI35hjJGJSZUnVmJ", + "kjtaHc6ThORhSm4PInkySVaMb7FM5klV6Z87ZL7CLkynxHKDJOs5x4FNNQ8rp6EdMvDlt5JxeQESk0J0", + "uX/Eu4LhHK0YR6CXIsnJeg0ccRAlowKmyeRAuXcsXcZUcsdSRHLEocAaB1YqQ3s6REuskmUllwXLsKEb", + "YuO+IkLR44ZkmwYX9Ccp0YoUgB5JUaAU0IpVNFfM4RveloXm92Y+mz0Z9e1nTwZ6JN/Pnu5Yqv+VZAtC", + "4m25Xz4ZwiTfT/8kZejQQmJZaeX8k8MqmSf/mNUOP7PePjPGuDFrD+1pteppxa1545kFoMSEIGnhdaHs", + "oQiC0PLTaqt4qbMVIEGxI3RZcrbmIJS/rjApIG8wr2U0zG8NPI5DyWJIxwPaMH4bRysCRb5ckUICD4jz", + "Tn/QhhUZKwGlWECOGEV6IzJAQQ+4qIx0RMK21wbv1F5DOqljCuYc77Sr0zUIdYDnnMtvRiXmeAtqf1ds", + "oPnSxdsAsrGQSH1GbFUTnEZAx2WUFBFoRfhAYqHoOsDNmUkD9xXwnQoYVSGFcvnU4Q/yAT4fdoQQ/pvG", + "i9hFYw5lrKi2tKP6jNGcuKgyECXnfs9+YiEbT2Hqi9OKQSmhSj1xIxr4hqmZb6gSsKoKF0oNLlt47xBt", + "o/pAvw0ZJg2N+KP0KP68qcMjAciIv/JG8ZxakagAIZZygxX/ta6BuPsxK5iAfEmoBP6Ai2SSsBJo8+cC", + "VnLZXcbJehP6PaFZUeWg84D5XyjMLfJ/gRB4HaiHQt7wmZL7ChDJVf5YEeAml1JJ5O507Edwf8XSU6qz", + "mxIydRaBHjcgNzZYaZ8VtkArdrpEU2cIFWkmH0bd39V9uSecQsHoWjn/kBx/PKwMoTAs4V6xNJJtNdVa", + "ykmtzaMJ+IqlpmT+TsVyrwGG6Op17hCatosaTtaYjlhRQCZZIEo3v6Ity6HoZkijluUDcBGNM1Z1do2L", + "uZkj3XI9QuWbn2sBVUhYm/Svg1Qk+ppvKoA7kLeoD6o2mtLqGBqqOY4l2sWFEw0LwTKCJeTqLINK6ULV", + "ffK4Hs2al+pxSyjZVttlpoAL/DhLu9jz1BVj8zplqZzIW8I3+TzGaqcq2cRAnrHK3aIpro7IYTuWmhy6", + "QJ+fGWzFYYxqB/NSdxA9tPNgC29DWC0LFja6KIzi2pSMOoFgaS5qJSsrpQd9OFNMnp5AbYDyzPsUV7c+", + "4reZSq0xdWVz698+bv3tnD6IpEjjbtQK6gD0Na8I3GMdg7pQ1LdZvUpBErsU48r0xYflzX8/nCeT5MP1", + "rf7v5UWwer5i6endxJfWqyHTvCdCuv6gON6+a/rKYA/xvceAYyje/66AW8HbarivP7SPpHYp5dsFg0+i", + "OO16b32ObQghhkLnpM8PZzr8DvTkjNEVWXepn+vfV9z0+ZzXesqD7hThy1n4hK9fRJ2iBWeqZS883Eqb", + "QhcXbaT0ynMIVvPzAIjdqoXBW4Em0alZOqqLQi8aQ0+CBoXHI/D4f1Ww3h3V3C0I+Ql0g6mrPhffgkmo", + "vlj4DCwZ0jDGa0yoGNatN4VW1Dn+E/AIxacSYOwGQlVWz6m6m9IdnqNHX2aw0VWYboB1hfika9EDGcAo", + "6HhiNgTjp7HIOJaU1c5OQHFp+erm+sPy8vfbT2/Pb68/JZPk/PrD7eXvt8t3n9+/D+ZnzTeWofsCvtmX", + "HwT9F8Z5RzQ7Md73O7SjXIzi2N2iY6/nByvWPcnbjwvvZm35tD3QNc82IKSVXQB/IJmpACSRel4VWvf2", + "4yKZJN5ayU/Ts+mZHqGVQHFJknnyZno2faPuKVhutJCzrG5fMRGYkJpYLBDW8dSWqY9EmglbydkDySFH", + "uZkoTk1X1hxokfv9tklmrARC/srynY3o0rLHZVkQM8Wb3QkDNxMO+4Jla2K+b/uc5BXoXxgf1zL/fPbT", + "q/Guu8Oa8YHujLYyfbIciSrLQIhVVRS7qbLLL2dngUBMH3ChR6ZaUyhluVq9nySi2m6xKsysUls2UdjA", + "a6GigGF7Y1CTfFFbrZlnTyTfK5br0DD8E0hO4EEbW5igkzmLpztEpECLi66FfwPZMG9Lz2evpue60RrX", + "s8fgUNU2p3J6zy8B/NuWKJNubNw2xW8gEW5oaXFxxBLK8RzLZP7H8QafpSkZ4tYwOjyphcp9Xfd1bmqu", + "NuQnDbX23R2/6HCQbWIBU5iZLdGZuRUBFPqijm82fwfHt9lrkOOfRQ1c2eTwIicdiKEDCBkBPIp6PPnO", + "zn4GRGt1Yz45VF+xdCRz1d2RHylIX7F07Ah9p1XqjHrF0o5FZ2Y8HjfsghJJtGkbD1tKztRxu5a0jzLM", + "C42RrNl+/vEjWfSy9XBpJJ9WsIkmBauV+rFLw/z2zU4HAaenZt0lPJqXnS+Pl5RjHjVKLj6qc5OIrU5a", + "WbjlcANSsKJhFPpXJ1sfwAdk2lHj9ik5NmK4MRPsUWD49DogDCsnnGXNafVp7ui39jumn6f8urti6SIf", + "u3r2QsVyoD9623cj6q6XD/DIlloOfdOf7BlO6u/N7ZHWD+C2tcgnOHAQE+MUYgcDz5d4di3q2D4+AHRN", + "b6/X9yKuFQDy5pyoJwAo8LmGUj1ASnd+ZBWKAK1h1Hfw//bwK5ay/emVPM8MtrppLlvDNOWigeDrz/MS", + "t/dcRnL7FizqGj2IifMNZF8NILKKc3Xj0y+t9QC1+eq5jYX6DTeMCYH2w/941bzBAqUAFPk34XEkXNZv", + "yqOA0FpxFxatj+mRYvikBkmt0ZEsf+/moL1FQNGe1TafpuvH4RknEjjB4VDgZsMjRwDHJmD8t4cSvG7h", + "brrVx4PGvVeCA4fe1cLGgJ6H6eWf3PUwI+9x0m1zpvkj3ZONVUbufdxbxUaM6t3sGfdfY+rjhXZt15Ec", + "y762iGl3lHtwjzuZuttrp1Vvd31qeLz1s7e/ur5u+PiAunp0135hA9pYc8zSuQcvvmQ+yVdn0j0y+GEQ", + "FEwPl98gq0x+0K8ZeEXNk7ZYpmj/rdZBYxWEHB1PzUccw0H12vzto4hAZLvVarQPcl+SKDShftipPcAf", + "wgD7yFleZX5ODuqqV/EimScbKUsxn81wSab2Dz2nGdvOHn5KFFost0N6187ewvzRKuRIMjuZETVa25OZ", + "/WQYmTuWNmg0OlAnEKgvtW1SnUvtKTTrG2+TZOeKNpRkXctZYi2TDqVi6vsGlfZtYf9l/78AAAD//yRV", + "PgtHPwAA", } // GetSwagger returns the content of the embedded swagger specification file diff --git a/api/queryService/controllers.go b/api/queryService/controllers.go index a7b38274..970146ab 100644 --- a/api/queryService/controllers.go +++ b/api/queryService/controllers.go @@ -2,6 +2,7 @@ package queryservice import ( "queryorchestration/internal/client" + "queryorchestration/internal/document" "queryorchestration/internal/export" "queryorchestration/internal/job" "queryorchestration/internal/job/collector" @@ -24,6 +25,7 @@ type Services struct { QueryTest *querytest.Service Client *client.Service Job *job.Service + Document *document.Service } type Controllers struct { diff --git a/api/queryService/documents.go b/api/queryService/documents.go new file mode 100644 index 00000000..ffbae00d --- /dev/null +++ b/api/queryService/documents.go @@ -0,0 +1,27 @@ +package queryservice + +import ( + "fmt" + "net/http" + + "github.com/labstack/echo/v4" + "github.com/oapi-codegen/runtime/types" +) + +func (s *Controllers) ListDocumentsByJobId(ctx echo.Context, jobId types.UUID) error { + documents, err := s.svc.Document.ListByJobId(ctx.Request().Context(), jobId) + if err != nil { + return echo.NewHTTPError(http.StatusNotFound, fmt.Sprintf("Unable to list documents: %s", err)) + } + + docs := make([]Document, len(documents)) + for i, doc := range documents { + docs[i] = Document{ + Id: doc.ID, + Bucket: doc.Bucket, + Key: doc.Key, + } + } + + return ctx.JSON(http.StatusOK, docs) +} diff --git a/api/queryService/documents_test.go b/api/queryService/documents_test.go new file mode 100644 index 00000000..5bcd1bee --- /dev/null +++ b/api/queryService/documents_test.go @@ -0,0 +1,73 @@ +package queryservice_test + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + queryservice "queryorchestration/api/queryService" + "queryorchestration/internal/database" + "queryorchestration/internal/database/repository" + "queryorchestration/internal/document" + "queryorchestration/internal/serviceconfig" + "queryorchestration/internal/serviceconfig/queue/jobsync" + queuemock "queryorchestration/mocks/queue" + "strings" + "testing" + + "github.com/go-playground/validator/v10" + "github.com/google/uuid" + "github.com/labstack/echo/v4" + "github.com/pashagolub/pgxmock/v3" + "github.com/stretchr/testify/assert" +) + +func TestListDocumentsByJobId(t *testing.T) { + pool, err := pgxmock.NewPool() + if err != nil { + t.Fatalf("failed to open pgxmock database: %v", err) + } + cfg := &struct { + serviceconfig.BaseConfig + jobsync.JobSyncConfig + }{} + cfg.DBPool = pool + cfg.DBQueries = repository.New(pool) + mockSQS := queuemock.NewMockSQSClient(t) + cfg.QueueClient = mockSQS + + jobId := uuid.New() + + e := echo.New() + req := httptest.NewRequest(http.MethodPost, "/", strings.NewReader("")) + req.Header.Set(echo.HeaderContentType, echo.MIMEApplicationJSON) + rec := httptest.NewRecorder() + ctx := e.NewContext(req, rec) + + cons := queryservice.NewControllers(validator.New(), &queryservice.Services{ + Document: document.New(cfg), + }) + + ctx.Set("id", jobId) + doc := queryservice.ListDocuments{ + { + Id: uuid.New(), + Bucket: "bucketone", + Key: "keyone", + }, + } + + pool.ExpectQuery("name: ListDocumentsByJobId :many").WithArgs(database.MustToDBUUID(jobId)). + WillReturnRows( + pgxmock.NewRows([]string{"id", "bucket", "key"}). + AddRow(database.MustToDBUUID(doc[0].Id), doc[0].Bucket, doc[0].Key), + ) + + err = cons.ListDocumentsByJobId(ctx, jobId) + assert.NoError(t, err) + assert.Equal(t, http.StatusOK, rec.Code) + + var res queryservice.ListDocuments + err = json.Unmarshal(rec.Body.Bytes(), &res) + assert.NoError(t, err) + assert.EqualExportedValues(t, doc, res) +} diff --git a/cmd/queryService/main.go b/cmd/queryService/main.go index 23d7291b..d707e715 100644 --- a/cmd/queryService/main.go +++ b/cmd/queryService/main.go @@ -78,6 +78,7 @@ func main() { QueryTest: quetest, Client: cli, Job: jbb, + Document: doc, } cons := queryservice.NewControllers(cfg.GetValidator(), services) diff --git a/database/queries/document.sql b/database/queries/document.sql index 2db398f6..d81149e8 100644 --- a/database/queries/document.sql +++ b/database/queries/document.sql @@ -1,6 +1,21 @@ -- name: GetDocument :one SELECT id, jobId, hash FROM documents WHERE id = $1 LIMIT 1; +-- name: ListDocumentsByJobId :many +WITH docs as ( + SELECT id from documents where jobId = @jobId +), +entries as ( + SELECT + de.documentId, + de.bucket, + de.key, + ROW_NUMBER() OVER (PARTITION BY de.documentId ORDER BY de.id DESC) as rowNumber + FROM docs d + JOIN documentEntries de on d.id = de.documentId +) +SELECT documentId as id, bucket, key from entries WHERE rowNumber = 1; + -- name: CreateDocument :one INSERT INTO documents (jobId, hash) VALUES ($1, $2) RETURNING id; diff --git a/internal/database/repository/document.sql.go b/internal/database/repository/document.sql.go index 32410885..40cf092f 100644 --- a/internal/database/repository/document.sql.go +++ b/internal/database/repository/document.sql.go @@ -100,3 +100,60 @@ func (q *Queries) GetDocumentIDByHash(ctx context.Context, arg *GetDocumentIDByH err := row.Scan(&id) return id, err } + +const listDocumentsByJobId = `-- name: ListDocumentsByJobId :many +WITH docs as ( + SELECT id from documents where jobId = $1 +), +entries as ( + SELECT + de.documentId, + de.bucket, + de.key, + ROW_NUMBER() OVER (PARTITION BY de.documentId ORDER BY de.id DESC) as rowNumber + FROM docs d + JOIN documentEntries de on d.id = de.documentId +) +SELECT documentId as id, bucket, key from entries WHERE rowNumber = 1 +` + +type ListDocumentsByJobIdRow struct { + ID pgtype.UUID `db:"id"` + Bucket string `db:"bucket"` + Key string `db:"key"` +} + +// ListDocumentsByJobId +// +// WITH docs as ( +// SELECT id from documents where jobId = $1 +// ), +// entries as ( +// SELECT +// de.documentId, +// de.bucket, +// de.key, +// ROW_NUMBER() OVER (PARTITION BY de.documentId ORDER BY de.id DESC) as rowNumber +// FROM docs d +// JOIN documentEntries de on d.id = de.documentId +// ) +// SELECT documentId as id, bucket, key from entries WHERE rowNumber = 1 +func (q *Queries) ListDocumentsByJobId(ctx context.Context, jobid pgtype.UUID) ([]*ListDocumentsByJobIdRow, error) { + rows, err := q.db.Query(ctx, listDocumentsByJobId, jobid) + if err != nil { + return nil, err + } + defer rows.Close() + items := []*ListDocumentsByJobIdRow{} + for rows.Next() { + var i ListDocumentsByJobIdRow + if err := rows.Scan(&i.ID, &i.Bucket, &i.Key); err != nil { + return nil, err + } + items = append(items, &i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} diff --git a/internal/database/repository/document_test.go b/internal/database/repository/document_test.go index 2accd499..4d168fac 100644 --- a/internal/database/repository/document_test.go +++ b/internal/database/repository/document_test.go @@ -34,6 +34,10 @@ func TestDocument(t *testing.T) { jobId, err := queries.CreateJob(ctx, clientId) assert.NoError(t, err) + docs, err := queries.ListDocumentsByJobId(ctx, jobId) + assert.NoError(t, err) + assert.Len(t, docs, 0) + hash := "example_hash" id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Jobid: jobId, @@ -41,20 +45,76 @@ func TestDocument(t *testing.T) { }) assert.NoError(t, err) assert.NotEmpty(t, id) + + bucketone := "bucketone" + keyone := "keyone" + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: id, + Bucket: bucketone, + Key: keyone, + }) + assert.NoError(t, err) + + docs, err = queries.ListDocumentsByJobId(ctx, jobId) + assert.NoError(t, err) + assert.Len(t, docs, 1) + assert.Equal(t, bucketone, docs[0].Bucket) + assert.Equal(t, keyone, docs[0].Key) + documentTwoID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Jobid: jobId, Hash: "example_hash_two", }) assert.NoError(t, err) assert.NotEmpty(t, documentTwoID) + + buckettwo := "buckettwo" + keytwo := "keytwo" + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: documentTwoID, + Bucket: buckettwo, + Key: keytwo, + }) + assert.NoError(t, err) + + docs, err = queries.ListDocumentsByJobId(ctx, jobId) + assert.NoError(t, err) + assert.Len(t, docs, 2) + jobTwoId, err := queries.CreateJob(ctx, clientId) assert.NoError(t, err) - _, err = queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + + docs, err = queries.ListDocumentsByJobId(ctx, jobId) + assert.NoError(t, err) + assert.Len(t, docs, 2) + docs, err = queries.ListDocumentsByJobId(ctx, jobTwoId) + assert.NoError(t, err) + assert.Len(t, docs, 0) + + documentThreeId, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Jobid: jobTwoId, Hash: "example_hash", }) assert.NoError(t, err) + bucketthree := "buckettwo" + keythree := "keytwo" + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: documentThreeId, + Bucket: bucketthree, + Key: keythree, + }) + assert.NoError(t, err) + + docs, err = queries.ListDocumentsByJobId(ctx, jobId) + assert.NoError(t, err) + assert.Len(t, docs, 2) + docs, err = queries.ListDocumentsByJobId(ctx, jobTwoId) + assert.NoError(t, err) + assert.Len(t, docs, 1) + assert.Equal(t, bucketthree, docs[0].Bucket) + assert.Equal(t, keythree, docs[0].Key) + doc, err := queries.GetDocument(ctx, id) assert.NoError(t, err) assert.EqualExportedValues(t, &repository.Document{ @@ -70,48 +130,27 @@ func TestDocument(t *testing.T) { assert.NoError(t, err) assert.EqualExportedValues(t, id, docid) - err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ - Documentid: id, - Bucket: "bucket_one", - Key: "/i/am/here", - }) - assert.NoError(t, err) - entry, err := queries.GetDocumentEntry(ctx, id) assert.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ Documentid: id, - Bucket: "bucket_one", - Key: "/i/am/here", + Bucket: bucketone, + Key: keyone, }, entry) - err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ - Documentid: id, - Bucket: "bucket_two", - Key: "/you/is/there", - }) - assert.NoError(t, err) - - entry, err = queries.GetDocumentEntry(ctx, id) + entry, err = queries.GetDocumentEntry(ctx, documentTwoID) assert.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ - Documentid: id, - Bucket: "bucket_two", - Key: "/you/is/there", + Documentid: documentTwoID, + Bucket: buckettwo, + Key: keytwo, }, entry) - err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ - Documentid: id, - Bucket: "bucket_three", - Key: "/who/is/where", - }) - assert.NoError(t, err) - - entry, err = queries.GetDocumentEntry(ctx, id) + entry, err = queries.GetDocumentEntry(ctx, documentThreeId) assert.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ - Documentid: id, - Bucket: "bucket_three", - Key: "/who/is/where", + Documentid: documentThreeId, + Bucket: bucketthree, + Key: keythree, }, entry) } diff --git a/internal/document/list.go b/internal/document/list.go new file mode 100644 index 00000000..02fd9eff --- /dev/null +++ b/internal/document/list.go @@ -0,0 +1,32 @@ +package document + +import ( + "context" + "queryorchestration/internal/database" + + "github.com/google/uuid" +) + +type DocumentInList struct { + ID uuid.UUID + Bucket string + Key string +} + +func (s *Service) ListByJobId(ctx context.Context, id uuid.UUID) ([]*DocumentInList, error) { + documents, err := s.cfg.GetDBQueries().ListDocumentsByJobId(ctx, database.MustToDBUUID(id)) + if err != nil { + return nil, err + } + + docs := make([]*DocumentInList, len(documents)) + for i, doc := range documents { + docs[i] = &DocumentInList{ + ID: database.MustToUUID(doc.ID), + Bucket: doc.Bucket, + Key: doc.Key, + } + } + + return docs, nil +} diff --git a/internal/document/list_test.go b/internal/document/list_test.go new file mode 100644 index 00000000..91c6bd84 --- /dev/null +++ b/internal/document/list_test.go @@ -0,0 +1,47 @@ +package document_test + +import ( + "context" + "queryorchestration/internal/database" + "queryorchestration/internal/database/repository" + "queryorchestration/internal/document" + "queryorchestration/internal/serviceconfig" + "testing" + + "github.com/google/uuid" + "github.com/pashagolub/pgxmock/v3" + "github.com/stretchr/testify/assert" +) + +func TestListByJobId(t *testing.T) { + ctx := context.Background() + + pool, err := pgxmock.NewPool() + if err != nil { + t.Fatalf("failed to open pgxmock database: %v", err) + } + cfg := &serviceconfig.BaseConfig{} + cfg.DBPool = pool + cfg.DBQueries = repository.New(pool) + + svc := document.New(cfg) + + jobId := uuid.New() + doc := []*document.DocumentInList{ + { + ID: uuid.New(), + Bucket: "bucketone", + Key: "keyone", + }, + } + + pool.ExpectQuery("name: ListDocumentsByJobId :many").WithArgs(database.MustToDBUUID(jobId)). + WillReturnRows( + pgxmock.NewRows([]string{"id", "bucket", "key"}). + AddRow(database.MustToDBUUID(doc[0].ID), doc[0].Bucket, doc[0].Key), + ) + + adoc, err := svc.ListByJobId(ctx, jobId) + assert.NoError(t, err) + assert.Equal(t, doc, adoc) +} diff --git a/internal/query/types/contextFull/process_test.go b/internal/query/types/contextFull/process_test.go index 2dc3e1bc..ccf94c18 100644 --- a/internal/query/types/contextFull/process_test.go +++ b/internal/query/types/contextFull/process_test.go @@ -25,7 +25,7 @@ func TestContextFull(t *testing.T) { value, err := extractor.Process(ctx, query, values) assert.NoError(t, err) - assert.Equal(t, "{\"id\":\"aaaa\",\"two\":\"bbbb\",\"three\":\"ccc\"}", value) + assert.Equal(t, `{"keyone":"valueone","keytwo":"valuetwo"}`, value) values = []resultprocessor.Value{ contextfull.NewResult("example_result"), diff --git a/internal/query/types/contextFull/service.go b/internal/query/types/contextFull/service.go index 99df26fe..16ef84f2 100644 --- a/internal/query/types/contextFull/service.go +++ b/internal/query/types/contextFull/service.go @@ -19,5 +19,5 @@ func (e *Extractor) Process(ctx context.Context, query *resultprocessor.Query, v } // TODO - return `{"id":"aaaa","two":"bbbb","three":"ccc"}`, nil + return `{"keyone":"valueone","keytwo":"valuetwo"}`, nil } diff --git a/pkg/queryService/api.gen.go b/pkg/queryService/api.gen.go index a05dd30c..f547ddcb 100644 --- a/pkg/queryService/api.gen.go +++ b/pkg/queryService/api.gen.go @@ -63,6 +63,18 @@ type ClientUpdate struct { Name *string `json:"name,omitempty"` } +// Document defines model for Document. +type Document struct { + // Bucket The bucket containing the document + Bucket string `json:"bucket"` + + // Id The document id + Id openapi_types.UUID `json:"id"` + + // Key The path to the document + Key string `json:"key"` +} + // ExportDetails Payload for export trigger response. type ExportDetails struct { // JobId The job id relative to the export. @@ -204,6 +216,9 @@ type JobUpdate struct { CanSync *bool `json:"can_sync,omitempty"` } +// ListDocuments The documents in the job. +type ListDocuments = []Document + // ListQueries defines model for ListQueries. type ListQueries struct { // Queries List of queries. @@ -412,6 +427,9 @@ type ClientInterface interface { UpdateJobCollectorByJobId(ctx context.Context, id openapi_types.UUID, body UpdateJobCollectorByJobIdJSONRequestBody, reqEditors ...RequestEditorFn) (*http.Response, error) + // ListDocumentsByJobId request + ListDocumentsByJobId(ctx context.Context, id openapi_types.UUID, reqEditors ...RequestEditorFn) (*http.Response, error) + // ExportState request ExportState(ctx context.Context, id openapi_types.UUID, reqEditors ...RequestEditorFn) (*http.Response, error) @@ -617,6 +635,18 @@ func (c *Client) UpdateJobCollectorByJobId(ctx context.Context, id openapi_types return c.Client.Do(req) } +func (c *Client) ListDocumentsByJobId(ctx context.Context, id openapi_types.UUID, reqEditors ...RequestEditorFn) (*http.Response, error) { + req, err := NewListDocumentsByJobIdRequest(c.Server, id) + if err != nil { + return nil, err + } + req = req.WithContext(ctx) + if err := c.applyEditors(ctx, req, reqEditors); err != nil { + return nil, err + } + return c.Client.Do(req) +} + func (c *Client) ExportState(ctx context.Context, id openapi_types.UUID, reqEditors ...RequestEditorFn) (*http.Response, error) { req, err := NewExportStateRequest(c.Server, id) if err != nil { @@ -1088,6 +1118,40 @@ func NewUpdateJobCollectorByJobIdRequestWithBody(server string, id openapi_types return req, nil } +// NewListDocumentsByJobIdRequest generates requests for ListDocumentsByJobId +func NewListDocumentsByJobIdRequest(server string, id openapi_types.UUID) (*http.Request, error) { + var err error + + var pathParam0 string + + pathParam0, err = runtime.StyleParamWithLocation("simple", false, "id", runtime.ParamLocationPath, id) + if err != nil { + return nil, err + } + + serverURL, err := url.Parse(server) + if err != nil { + return nil, err + } + + operationPath := fmt.Sprintf("/job/%s/documents", pathParam0) + if operationPath[0] == '/' { + operationPath = "." + operationPath + } + + queryURL, err := serverURL.Parse(operationPath) + if err != nil { + return nil, err + } + + req, err := http.NewRequest("GET", queryURL.String(), nil) + if err != nil { + return nil, err + } + + return req, nil +} + // NewExportStateRequest generates requests for ExportState func NewExportStateRequest(server string, id openapi_types.UUID) (*http.Request, error) { var err error @@ -1399,6 +1463,9 @@ type ClientWithResponsesInterface interface { UpdateJobCollectorByJobIdWithResponse(ctx context.Context, id openapi_types.UUID, body UpdateJobCollectorByJobIdJSONRequestBody, reqEditors ...RequestEditorFn) (*UpdateJobCollectorByJobIdResponse, error) + // ListDocumentsByJobIdWithResponse request + ListDocumentsByJobIdWithResponse(ctx context.Context, id openapi_types.UUID, reqEditors ...RequestEditorFn) (*ListDocumentsByJobIdResponse, error) + // ExportStateWithResponse request ExportStateWithResponse(ctx context.Context, id openapi_types.UUID, reqEditors ...RequestEditorFn) (*ExportStateResponse, error) @@ -1619,6 +1686,28 @@ func (r UpdateJobCollectorByJobIdResponse) StatusCode() int { return 0 } +type ListDocumentsByJobIdResponse struct { + Body []byte + HTTPResponse *http.Response + JSON200 *ListDocuments +} + +// Status returns HTTPResponse.Status +func (r ListDocumentsByJobIdResponse) Status() string { + if r.HTTPResponse != nil { + return r.HTTPResponse.Status + } + return http.StatusText(0) +} + +// StatusCode returns HTTPResponse.StatusCode +func (r ListDocumentsByJobIdResponse) StatusCode() int { + if r.HTTPResponse != nil { + return r.HTTPResponse.StatusCode + } + return 0 +} + type ExportStateResponse struct { Body []byte HTTPResponse *http.Response @@ -1879,6 +1968,15 @@ func (c *ClientWithResponses) UpdateJobCollectorByJobIdWithResponse(ctx context. return ParseUpdateJobCollectorByJobIdResponse(rsp) } +// ListDocumentsByJobIdWithResponse request returning *ListDocumentsByJobIdResponse +func (c *ClientWithResponses) ListDocumentsByJobIdWithResponse(ctx context.Context, id openapi_types.UUID, reqEditors ...RequestEditorFn) (*ListDocumentsByJobIdResponse, error) { + rsp, err := c.ListDocumentsByJobId(ctx, id, reqEditors...) + if err != nil { + return nil, err + } + return ParseListDocumentsByJobIdResponse(rsp) +} + // ExportStateWithResponse request returning *ExportStateResponse func (c *ClientWithResponses) ExportStateWithResponse(ctx context.Context, id openapi_types.UUID, reqEditors ...RequestEditorFn) (*ExportStateResponse, error) { rsp, err := c.ExportState(ctx, id, reqEditors...) @@ -2161,6 +2259,32 @@ func ParseUpdateJobCollectorByJobIdResponse(rsp *http.Response) (*UpdateJobColle return response, nil } +// ParseListDocumentsByJobIdResponse parses an HTTP response from a ListDocumentsByJobIdWithResponse call +func ParseListDocumentsByJobIdResponse(rsp *http.Response) (*ListDocumentsByJobIdResponse, error) { + bodyBytes, err := io.ReadAll(rsp.Body) + defer func() { _ = rsp.Body.Close() }() + if err != nil { + return nil, err + } + + response := &ListDocumentsByJobIdResponse{ + Body: bodyBytes, + HTTPResponse: rsp, + } + + switch { + case strings.Contains(rsp.Header.Get("Content-Type"), "json") && rsp.StatusCode == 200: + var dest ListDocuments + if err := json.Unmarshal(bodyBytes, &dest); err != nil { + return nil, err + } + response.JSON200 = &dest + + } + + return response, nil +} + // ParseExportStateResponse parses an HTTP response from a ExportStateWithResponse call func ParseExportStateResponse(rsp *http.Response) (*ExportStateResponse, error) { bodyBytes, err := io.ReadAll(rsp.Body) diff --git a/serviceAPIs/queryService.yaml b/serviceAPIs/queryService.yaml index f14faa98..1c6cf74e 100644 --- a/serviceAPIs/queryService.yaml +++ b/serviceAPIs/queryService.yaml @@ -14,6 +14,8 @@ tags: description: Operations related to jobs - name: JobCollectorService description: Operations related to job collectors + - name: JobDocumentsService + description: Operations related to job documents - name: QueryService description: Operations related to queries - name: ExportService @@ -319,6 +321,31 @@ paths: "404": description: Job collector not found. + /job/{id}/documents: + parameters: + - in: path + name: id + required: true + schema: + type: string + format: uuid + description: The job ID for the documents. + get: + operationId: listDocumentsByJobId + tags: + - JobDocumentsService + summary: List the documents for a job + description: Retrieves the list of documents by the job ID. + responses: + "200": + description: Job documents list. + content: + application/json: + schema: + $ref: "#/components/schemas/ListDocuments" + "404": + description: Job not found. + /job/export: post: operationId: triggerExport @@ -543,6 +570,30 @@ components: type: boolean description: Specifies whether the job is actively syncing + ListDocuments: + type: array + description: The documents in the job. + items: + $ref: "#/components/schemas/Document" + + Document: + type: object + properties: + id: + type: string + format: uuid + description: The document id + bucket: + type: string + description: The bucket containing the document + key: + type: string + description: The path to the document + required: + - id + - bucket + - key + ClientCreate: type: object properties: diff --git a/test/process_test.go b/test/process_test.go index 19e09d21..8796e94a 100644 --- a/test/process_test.go +++ b/test/process_test.go @@ -12,6 +12,7 @@ import ( queryservice "queryorchestration/pkg/queryService" "strings" "testing" + "time" "github.com/aws/aws-sdk-go-v2/service/s3" awstypes "github.com/aws/aws-sdk-go-v2/service/s3/types" @@ -172,7 +173,7 @@ func TestProcess(t *testing.T) { Type: queryservice.CONTEXTFULL, }) assert.NoError(t, err) - jcfg := `{"path":"id"}` + jcfg := `{"path":"keyone"}` jsonQueryRes, err := qService.CreateQueryWithResponse(ctx, queryservice.QueryCreate{ Type: queryservice.JSONEXTRACTOR, Config: &jcfg, @@ -191,12 +192,75 @@ func TestProcess(t *testing.T) { }) assert.NoError(t, err) + jRes, err := qService.GetJobWithResponse(ctx, jobRes.JSON201.Id) + assert.NoError(t, err) + assert.Equal(t, queryservice.INSYNC, jRes.JSON200.Status) + location := fmt.Sprintf("%s/%s/%s", clientRes.JSON201.Id, jobRes.JSON201.Id, "object_name") - body := strings.NewReader("hello world") + body := strings.NewReader(`{"keyone":"valueone","keytwo":"valuetwo"}`) _, err = cfg.StoreClient.PutObject(ctx, &s3.PutObjectInput{ Bucket: &bucketName, Key: &location, Body: body, }) assert.NoError(t, err) + + WaitForJobStatus(t, ctx, qService, jobRes.JSON201.Id, queryservice.NOTSYNCED) + WaitForJobStatus(t, ctx, qService, jobRes.JSON201.Id, queryservice.INSYNC) + + docs, err := qService.ListDocumentsByJobIdWithResponse(ctx, jobRes.JSON201.Id) + assert.NoError(t, err) + assert.Len(t, *docs.JSON200, 1) + doc := (*docs.JSON200)[0] + assert.Equal(t, bucketName, doc.Bucket) + assert.Equal(t, location, doc.Key) + + testRes, err := qService.TestQueryWithResponse(ctx, jsonQueryRes.JSON201.Id, queryservice.QueryTestRequest{ + QueryVersion: 1, + DocumentId: doc.Id, + }) + assert.NoError(t, err) + assert.Equal(t, "valueone", testRes.JSON200.Value) + + aV := int32(2) + jcfg = `{"path":"keytwo"}` + res, err := qService.UpdateQueryWithResponse(ctx, jsonQueryRes.JSON201.Id, queryservice.QueryUpdate{ + ActiveVersion: &aV, + Config: &jcfg, + }) + assert.NoError(t, err) + assert.NotNil(t, res) + + WaitForJobStatus(t, ctx, qService, jobRes.JSON201.Id, queryservice.NOTSYNCED) + WaitForJobStatus(t, ctx, qService, jobRes.JSON201.Id, queryservice.INSYNC) + + testRes, err = qService.TestQueryWithResponse(ctx, jsonQueryRes.JSON201.Id, queryservice.QueryTestRequest{ + QueryVersion: 2, + DocumentId: doc.Id, + }) + assert.NoError(t, err) + assert.Equal(t, "valuetwo", testRes.JSON200.Value) +} + +func WaitForJobStatus(t testing.TB, ctx context.Context, service *queryservice.ClientWithResponses, id types.UUID, status queryservice.JobStatus) { + timeout := time.After(30 * time.Second) + ticker := time.NewTicker(500 * time.Millisecond) + defer ticker.Stop() + + for { + select { + case <-timeout: + t.Fatalf("Timeout waiting for job status to become %s", status) + case <-ticker.C: + jRes, err := service.GetJobWithResponse(ctx, id) + if err != nil { + assert.NoError(t, err) + } + + if jRes.JSON200.Status == status { + assert.Equal(t, status, jRes.JSON200.Status) + return + } + } + } } diff --git a/test/queryService/openapi_test.go b/test/queryService/accessory_test.go similarity index 81% rename from test/queryService/openapi_test.go rename to test/queryService/accessory_test.go index 45eecc49..aec9a64d 100644 --- a/test/queryService/openapi_test.go +++ b/test/queryService/accessory_test.go @@ -10,7 +10,7 @@ import ( "github.com/stretchr/testify/assert" ) -func TestQueryServiceOpenAPI(t *testing.T) { +func TestQueryServiceAccessories(t *testing.T) { ctx := context.Background() c, cleanup := test.CreateServiceNetwork(t, ctx, &test.ServiceNetworkConfig{ @@ -36,4 +36,9 @@ func TestQueryServiceOpenAPI(t *testing.T) { assert.NoError(t, err) assert.NotNil(t, resp) assert.Equal(t, http.StatusOK, resp.StatusCode) + + resp, err = http.Get(fmt.Sprintf("%s/metrics", c.URI)) + assert.NoError(t, err) + assert.NotNil(t, resp) + assert.Equal(t, http.StatusOK, resp.StatusCode) }