Merged in feature/e2etest (pull request #74)

Feature/e2etest

* e2etest

* testcleanups
This commit is contained in:
Michael McGuinness
2025-02-21 17:05:37 +00:00
parent 3d434eedb8
commit 08d2296717
16 changed files with 654 additions and 79 deletions
+80 -42
View File
@@ -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
+2
View File
@@ -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 {
+27
View File
@@ -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)
}
+73
View File
@@ -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)
}
+1
View File
@@ -78,6 +78,7 @@ func main() {
QueryTest: quetest,
Client: cli,
Job: jbb,
Document: doc,
}
cons := queryservice.NewControllers(cfg.GetValidator(), services)
+15
View File
@@ -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;
@@ -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
}
+71 -32
View File
@@ -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)
}
+32
View File
@@ -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
}
+47
View File
@@ -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)
}
@@ -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"),
+1 -1
View File
@@ -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
}
+124
View File
@@ -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)
+51
View File
@@ -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:
+66 -2
View File
@@ -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
}
}
}
}
@@ -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)
}