mirror of
https://github.com/pikami/cosmium.git
synced 2026-01-26 21:02:58 +00:00
Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d27c633e1d | ||
|
|
3987df89c0 | ||
|
|
6e3f4169a1 | ||
|
|
14c5400d23 | ||
|
|
1cf5ae92f4 | ||
|
|
5d99b653cc | ||
|
|
787cdb33cf | ||
|
|
5caa829ac1 |
@@ -15,7 +15,7 @@ jobs:
|
|||||||
uses: crazy-max/ghaction-xgo@v3.1.0
|
uses: crazy-max/ghaction-xgo@v3.1.0
|
||||||
with:
|
with:
|
||||||
xgo_version: latest
|
xgo_version: latest
|
||||||
go_version: 1.22.0
|
go_version: 1.24.0
|
||||||
dest: dist
|
dest: dist
|
||||||
pkg: sharedlibrary
|
pkg: sharedlibrary
|
||||||
prefix: cosmium
|
prefix: cosmium
|
||||||
|
|||||||
4
.github/workflows/release.yml
vendored
4
.github/workflows/release.yml
vendored
@@ -21,13 +21,13 @@ jobs:
|
|||||||
- name: Set up Go
|
- name: Set up Go
|
||||||
uses: actions/setup-go@v5
|
uses: actions/setup-go@v5
|
||||||
with:
|
with:
|
||||||
go-version: 1.22.0
|
go-version: 1.24.0
|
||||||
|
|
||||||
- name: Cross-Compile with xgo
|
- name: Cross-Compile with xgo
|
||||||
uses: crazy-max/ghaction-xgo@v3.1.0
|
uses: crazy-max/ghaction-xgo@v3.1.0
|
||||||
with:
|
with:
|
||||||
xgo_version: latest
|
xgo_version: latest
|
||||||
go_version: 1.22.0
|
go_version: 1.24.0
|
||||||
dest: sharedlibrary_dist
|
dest: sharedlibrary_dist
|
||||||
pkg: sharedlibrary
|
pkg: sharedlibrary
|
||||||
prefix: cosmium
|
prefix: cosmium
|
||||||
|
|||||||
@@ -135,6 +135,12 @@ docker_manifests:
|
|||||||
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-amd64"
|
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-amd64"
|
||||||
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-arm64"
|
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-arm64"
|
||||||
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-arm64v8"
|
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-arm64v8"
|
||||||
|
- name_template: 'ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer'
|
||||||
|
skip_push: auto
|
||||||
|
image_templates:
|
||||||
|
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-amd64"
|
||||||
|
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-arm64"
|
||||||
|
- "ghcr.io/pikami/{{ .ProjectName }}:{{ .Version }}-explorer-arm64v8"
|
||||||
|
|
||||||
checksum:
|
checksum:
|
||||||
name_template: 'checksums.txt'
|
name_template: 'checksums.txt'
|
||||||
|
|||||||
2
Makefile
2
Makefile
@@ -9,7 +9,7 @@ SERVER_LOCATION=./cmd/server
|
|||||||
SHARED_LIB_LOCATION=./sharedlibrary
|
SHARED_LIB_LOCATION=./sharedlibrary
|
||||||
SHARED_LIB_OPT=-buildmode=c-shared
|
SHARED_LIB_OPT=-buildmode=c-shared
|
||||||
XGO_TARGETS=linux/amd64,linux/arm64,windows/amd64,windows/arm64,darwin/amd64,darwin/arm64
|
XGO_TARGETS=linux/amd64,linux/arm64,windows/amd64,windows/arm64,darwin/amd64,darwin/arm64
|
||||||
GOVERSION=1.22.0
|
GOVERSION=1.24.0
|
||||||
|
|
||||||
DIST_DIR=dist
|
DIST_DIR=dist
|
||||||
|
|
||||||
|
|||||||
24
api/api_models/models.go
Normal file
24
api/api_models/models.go
Normal file
@@ -0,0 +1,24 @@
|
|||||||
|
package apimodels
|
||||||
|
|
||||||
|
const (
|
||||||
|
BatchOperationTypeCreate = "Create"
|
||||||
|
BatchOperationTypeDelete = "Delete"
|
||||||
|
BatchOperationTypeReplace = "Replace"
|
||||||
|
BatchOperationTypeUpsert = "Upsert"
|
||||||
|
BatchOperationTypeRead = "Read"
|
||||||
|
BatchOperationTypePatch = "Patch"
|
||||||
|
)
|
||||||
|
|
||||||
|
type BatchOperation struct {
|
||||||
|
OperationType string `json:"operationType"`
|
||||||
|
Id string `json:"id"`
|
||||||
|
ResourceBody map[string]interface{} `json:"resourceBody"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type BatchOperationResult struct {
|
||||||
|
StatusCode int `json:"statusCode"`
|
||||||
|
RequestCharge float64 `json:"requestCharge"`
|
||||||
|
ResourceBody map[string]interface{} `json:"resourceBody"`
|
||||||
|
Etag string `json:"etag"`
|
||||||
|
Message string `json:"message"`
|
||||||
|
}
|
||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
|
|
||||||
jsonpatch "github.com/cosmiumdev/json-patch/v5"
|
jsonpatch "github.com/cosmiumdev/json-patch/v5"
|
||||||
"github.com/gin-gonic/gin"
|
"github.com/gin-gonic/gin"
|
||||||
|
apimodels "github.com/pikami/cosmium/api/api_models"
|
||||||
"github.com/pikami/cosmium/internal/constants"
|
"github.com/pikami/cosmium/internal/constants"
|
||||||
"github.com/pikami/cosmium/internal/logger"
|
"github.com/pikami/cosmium/internal/logger"
|
||||||
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
||||||
@@ -183,6 +184,13 @@ func (h *Handlers) DocumentsPost(c *gin.Context) {
|
|||||||
databaseId := c.Param("databaseId")
|
databaseId := c.Param("databaseId")
|
||||||
collectionId := c.Param("collId")
|
collectionId := c.Param("collId")
|
||||||
|
|
||||||
|
// Handle batch requests
|
||||||
|
isBatchRequest, _ := strconv.ParseBool(c.GetHeader("x-ms-cosmos-is-batch-request"))
|
||||||
|
if isBatchRequest {
|
||||||
|
h.handleBatchRequest(c)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
var requestBody map[string]interface{}
|
var requestBody map[string]interface{}
|
||||||
if err := c.BindJSON(&requestBody); err != nil {
|
if err := c.BindJSON(&requestBody); err != nil {
|
||||||
c.JSON(http.StatusBadRequest, gin.H{"message": err.Error()})
|
c.JSON(http.StatusBadRequest, gin.H{"message": err.Error()})
|
||||||
@@ -191,30 +199,7 @@ func (h *Handlers) DocumentsPost(c *gin.Context) {
|
|||||||
|
|
||||||
query := requestBody["query"]
|
query := requestBody["query"]
|
||||||
if query != nil {
|
if query != nil {
|
||||||
if c.GetHeader("x-ms-cosmos-is-query-plan-request") != "" {
|
h.handleDocumentQuery(c, requestBody)
|
||||||
c.IndentedJSON(http.StatusOK, constants.QueryPlanResponse)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
var queryParameters map[string]interface{}
|
|
||||||
if paramsArray, ok := requestBody["parameters"].([]interface{}); ok {
|
|
||||||
queryParameters = parametersToMap(paramsArray)
|
|
||||||
}
|
|
||||||
|
|
||||||
docs, status := h.repository.ExecuteQueryDocuments(databaseId, collectionId, query.(string), queryParameters)
|
|
||||||
if status != repositorymodels.StatusOk {
|
|
||||||
// TODO: Currently we return everything if the query fails
|
|
||||||
h.GetAllDocuments(c)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
collection, _ := h.repository.GetCollection(databaseId, collectionId)
|
|
||||||
c.Header("x-ms-item-count", fmt.Sprintf("%d", len(docs)))
|
|
||||||
c.IndentedJSON(http.StatusOK, gin.H{
|
|
||||||
"_rid": collection.ResourceID,
|
|
||||||
"Documents": docs,
|
|
||||||
"_count": len(docs),
|
|
||||||
})
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -253,3 +238,131 @@ func parametersToMap(pairs []interface{}) map[string]interface{} {
|
|||||||
|
|
||||||
return result
|
return result
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (h *Handlers) handleDocumentQuery(c *gin.Context, requestBody map[string]interface{}) {
|
||||||
|
databaseId := c.Param("databaseId")
|
||||||
|
collectionId := c.Param("collId")
|
||||||
|
|
||||||
|
if c.GetHeader("x-ms-cosmos-is-query-plan-request") != "" {
|
||||||
|
c.IndentedJSON(http.StatusOK, constants.QueryPlanResponse)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
var queryParameters map[string]interface{}
|
||||||
|
if paramsArray, ok := requestBody["parameters"].([]interface{}); ok {
|
||||||
|
queryParameters = parametersToMap(paramsArray)
|
||||||
|
}
|
||||||
|
|
||||||
|
docs, status := h.repository.ExecuteQueryDocuments(databaseId, collectionId, requestBody["query"].(string), queryParameters)
|
||||||
|
if status != repositorymodels.StatusOk {
|
||||||
|
// TODO: Currently we return everything if the query fails
|
||||||
|
h.GetAllDocuments(c)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
collection, _ := h.repository.GetCollection(databaseId, collectionId)
|
||||||
|
c.Header("x-ms-item-count", fmt.Sprintf("%d", len(docs)))
|
||||||
|
c.IndentedJSON(http.StatusOK, gin.H{
|
||||||
|
"_rid": collection.ResourceID,
|
||||||
|
"Documents": docs,
|
||||||
|
"_count": len(docs),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (h *Handlers) handleBatchRequest(c *gin.Context) {
|
||||||
|
databaseId := c.Param("databaseId")
|
||||||
|
collectionId := c.Param("collId")
|
||||||
|
|
||||||
|
batchOperations := make([]apimodels.BatchOperation, 0)
|
||||||
|
if err := c.BindJSON(&batchOperations); err != nil {
|
||||||
|
c.JSON(http.StatusBadRequest, gin.H{"message": err.Error()})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
batchOperationResults := make([]apimodels.BatchOperationResult, len(batchOperations))
|
||||||
|
for idx, operation := range batchOperations {
|
||||||
|
switch operation.OperationType {
|
||||||
|
case apimodels.BatchOperationTypeCreate:
|
||||||
|
createdDocument, status := h.repository.CreateDocument(databaseId, collectionId, operation.ResourceBody)
|
||||||
|
responseCode := repositoryStatusToResponseCode(status)
|
||||||
|
if status == repositorymodels.StatusOk {
|
||||||
|
responseCode = http.StatusCreated
|
||||||
|
}
|
||||||
|
batchOperationResults[idx] = apimodels.BatchOperationResult{
|
||||||
|
StatusCode: responseCode,
|
||||||
|
ResourceBody: createdDocument,
|
||||||
|
}
|
||||||
|
case apimodels.BatchOperationTypeDelete:
|
||||||
|
status := h.repository.DeleteDocument(databaseId, collectionId, operation.Id)
|
||||||
|
responseCode := repositoryStatusToResponseCode(status)
|
||||||
|
if status == repositorymodels.StatusOk {
|
||||||
|
responseCode = http.StatusNoContent
|
||||||
|
}
|
||||||
|
batchOperationResults[idx] = apimodels.BatchOperationResult{
|
||||||
|
StatusCode: responseCode,
|
||||||
|
}
|
||||||
|
case apimodels.BatchOperationTypeReplace:
|
||||||
|
deleteStatus := h.repository.DeleteDocument(databaseId, collectionId, operation.Id)
|
||||||
|
if deleteStatus == repositorymodels.StatusNotFound {
|
||||||
|
batchOperationResults[idx] = apimodels.BatchOperationResult{
|
||||||
|
StatusCode: http.StatusNotFound,
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
createdDocument, createStatus := h.repository.CreateDocument(databaseId, collectionId, operation.ResourceBody)
|
||||||
|
responseCode := repositoryStatusToResponseCode(createStatus)
|
||||||
|
if createStatus == repositorymodels.StatusOk {
|
||||||
|
responseCode = http.StatusCreated
|
||||||
|
}
|
||||||
|
batchOperationResults[idx] = apimodels.BatchOperationResult{
|
||||||
|
StatusCode: responseCode,
|
||||||
|
ResourceBody: createdDocument,
|
||||||
|
}
|
||||||
|
case apimodels.BatchOperationTypeUpsert:
|
||||||
|
documentId := operation.ResourceBody["id"].(string)
|
||||||
|
h.repository.DeleteDocument(databaseId, collectionId, documentId)
|
||||||
|
createdDocument, createStatus := h.repository.CreateDocument(databaseId, collectionId, operation.ResourceBody)
|
||||||
|
responseCode := repositoryStatusToResponseCode(createStatus)
|
||||||
|
if createStatus == repositorymodels.StatusOk {
|
||||||
|
responseCode = http.StatusCreated
|
||||||
|
}
|
||||||
|
batchOperationResults[idx] = apimodels.BatchOperationResult{
|
||||||
|
StatusCode: responseCode,
|
||||||
|
ResourceBody: createdDocument,
|
||||||
|
}
|
||||||
|
case apimodels.BatchOperationTypeRead:
|
||||||
|
document, status := h.repository.GetDocument(databaseId, collectionId, operation.Id)
|
||||||
|
batchOperationResults[idx] = apimodels.BatchOperationResult{
|
||||||
|
StatusCode: repositoryStatusToResponseCode(status),
|
||||||
|
ResourceBody: document,
|
||||||
|
}
|
||||||
|
case apimodels.BatchOperationTypePatch:
|
||||||
|
batchOperationResults[idx] = apimodels.BatchOperationResult{
|
||||||
|
StatusCode: http.StatusNotImplemented,
|
||||||
|
Message: "Patch operation is not implemented",
|
||||||
|
}
|
||||||
|
default:
|
||||||
|
batchOperationResults[idx] = apimodels.BatchOperationResult{
|
||||||
|
StatusCode: http.StatusBadRequest,
|
||||||
|
Message: "Unknown operation type",
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
c.JSON(http.StatusOK, batchOperationResults)
|
||||||
|
}
|
||||||
|
|
||||||
|
func repositoryStatusToResponseCode(status repositorymodels.RepositoryStatus) int {
|
||||||
|
switch status {
|
||||||
|
case repositorymodels.StatusOk:
|
||||||
|
return http.StatusOK
|
||||||
|
case repositorymodels.StatusNotFound:
|
||||||
|
return http.StatusNotFound
|
||||||
|
case repositorymodels.Conflict:
|
||||||
|
return http.StatusConflict
|
||||||
|
case repositorymodels.BadRequest:
|
||||||
|
return http.StatusBadRequest
|
||||||
|
default:
|
||||||
|
return http.StatusInternalServerError
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -75,8 +75,7 @@ func requestToResourceId(c *gin.Context) string {
|
|||||||
|
|
||||||
isFeed := c.Request.Header.Get("A-Im") == "Incremental Feed"
|
isFeed := c.Request.Header.Get("A-Im") == "Incremental Feed"
|
||||||
if resourceType == "pkranges" && isFeed {
|
if resourceType == "pkranges" && isFeed {
|
||||||
// CosmosSDK replaces '/' with '-' in resource id requests
|
resourceId = collId
|
||||||
resourceId = strings.Replace(collId, "-", "/", -1)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return resourceId
|
return resourceId
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import (
|
|||||||
|
|
||||||
"github.com/gin-gonic/gin"
|
"github.com/gin-gonic/gin"
|
||||||
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
||||||
|
"github.com/pikami/cosmium/internal/resourceid"
|
||||||
)
|
)
|
||||||
|
|
||||||
func (h *Handlers) GetPartitionKeyRanges(c *gin.Context) {
|
func (h *Handlers) GetPartitionKeyRanges(c *gin.Context) {
|
||||||
@@ -31,8 +32,9 @@ func (h *Handlers) GetPartitionKeyRanges(c *gin.Context) {
|
|||||||
collectionRid = collection.ResourceID
|
collectionRid = collection.ResourceID
|
||||||
}
|
}
|
||||||
|
|
||||||
|
rid := resourceid.NewCombined(collectionRid, resourceid.New(resourceid.ResourceTypePartitionKeyRange))
|
||||||
c.IndentedJSON(http.StatusOK, gin.H{
|
c.IndentedJSON(http.StatusOK, gin.H{
|
||||||
"_rid": collectionRid,
|
"_rid": rid,
|
||||||
"_count": len(partitionKeyRanges),
|
"_count": len(partitionKeyRanges),
|
||||||
"PartitionKeyRanges": partitionKeyRanges,
|
"PartitionKeyRanges": partitionKeyRanges,
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ import (
|
|||||||
|
|
||||||
"github.com/pikami/cosmium/api"
|
"github.com/pikami/cosmium/api"
|
||||||
"github.com/pikami/cosmium/api/config"
|
"github.com/pikami/cosmium/api/config"
|
||||||
|
"github.com/pikami/cosmium/internal/logger"
|
||||||
"github.com/pikami/cosmium/internal/repositories"
|
"github.com/pikami/cosmium/internal/repositories"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -35,6 +36,9 @@ func runTestServer() *TestServer {
|
|||||||
ExplorerBaseUrlLocation: config.ExplorerBaseUrlLocation,
|
ExplorerBaseUrlLocation: config.ExplorerBaseUrlLocation,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
config.LogLevel = "debug"
|
||||||
|
logger.SetLogLevel(logger.LogLevelDebug)
|
||||||
|
|
||||||
return runTestServerCustomConfig(config)
|
return runTestServerCustomConfig(config)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -377,5 +377,140 @@ func Test_Documents_Patch(t *testing.T) {
|
|||||||
assert.NotNil(t, r)
|
assert.NotNil(t, r)
|
||||||
assert.Nil(t, err2)
|
assert.Nil(t, err2)
|
||||||
})
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func Test_Documents_TransactionalBatch(t *testing.T) {
|
||||||
|
ts, collectionClient := documents_InitializeDb(t)
|
||||||
|
defer ts.Server.Close()
|
||||||
|
|
||||||
|
t.Run("Should execute CREATE transactional batch", func(t *testing.T) {
|
||||||
|
context := context.TODO()
|
||||||
|
batch := collectionClient.NewTransactionalBatch(azcosmos.NewPartitionKeyString("pk"))
|
||||||
|
|
||||||
|
newItem := map[string]interface{}{
|
||||||
|
"id": "678901",
|
||||||
|
}
|
||||||
|
bytes, err := json.Marshal(newItem)
|
||||||
|
assert.Nil(t, err)
|
||||||
|
|
||||||
|
batch.CreateItem(bytes, nil)
|
||||||
|
response, err := collectionClient.ExecuteTransactionalBatch(context, batch, &azcosmos.TransactionalBatchOptions{})
|
||||||
|
assert.Nil(t, err)
|
||||||
|
assert.True(t, response.Success)
|
||||||
|
assert.Equal(t, 1, len(response.OperationResults))
|
||||||
|
|
||||||
|
operationResponse := response.OperationResults[0]
|
||||||
|
assert.NotNil(t, operationResponse)
|
||||||
|
assert.NotNil(t, operationResponse.ResourceBody)
|
||||||
|
assert.Equal(t, int32(http.StatusCreated), operationResponse.StatusCode)
|
||||||
|
|
||||||
|
var itemResponseBody map[string]interface{}
|
||||||
|
json.Unmarshal(operationResponse.ResourceBody, &itemResponseBody)
|
||||||
|
assert.Equal(t, newItem["id"], itemResponseBody["id"])
|
||||||
|
|
||||||
|
createdDoc, _ := ts.Repository.GetDocument(testDatabaseName, testCollectionName, newItem["id"].(string))
|
||||||
|
assert.Equal(t, newItem["id"], createdDoc["id"])
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("Should execute DELETE transactional batch", func(t *testing.T) {
|
||||||
|
context := context.TODO()
|
||||||
|
batch := collectionClient.NewTransactionalBatch(azcosmos.NewPartitionKeyString("pk"))
|
||||||
|
|
||||||
|
batch.DeleteItem("12345", nil)
|
||||||
|
response, err := collectionClient.ExecuteTransactionalBatch(context, batch, &azcosmos.TransactionalBatchOptions{})
|
||||||
|
assert.Nil(t, err)
|
||||||
|
assert.True(t, response.Success)
|
||||||
|
assert.Equal(t, 1, len(response.OperationResults))
|
||||||
|
|
||||||
|
operationResponse := response.OperationResults[0]
|
||||||
|
assert.NotNil(t, operationResponse)
|
||||||
|
assert.Equal(t, int32(http.StatusNoContent), operationResponse.StatusCode)
|
||||||
|
|
||||||
|
_, status := ts.Repository.GetDocument(testDatabaseName, testCollectionName, "12345")
|
||||||
|
assert.Equal(t, repositorymodels.StatusNotFound, int(status))
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("Should execute REPLACE transactional batch", func(t *testing.T) {
|
||||||
|
context := context.TODO()
|
||||||
|
batch := collectionClient.NewTransactionalBatch(azcosmos.NewPartitionKeyString("pk"))
|
||||||
|
|
||||||
|
newItem := map[string]interface{}{
|
||||||
|
"id": "67890",
|
||||||
|
"pk": "666",
|
||||||
|
}
|
||||||
|
bytes, err := json.Marshal(newItem)
|
||||||
|
assert.Nil(t, err)
|
||||||
|
|
||||||
|
batch.ReplaceItem("67890", bytes, nil)
|
||||||
|
response, err := collectionClient.ExecuteTransactionalBatch(context, batch, &azcosmos.TransactionalBatchOptions{})
|
||||||
|
assert.Nil(t, err)
|
||||||
|
assert.True(t, response.Success)
|
||||||
|
assert.Equal(t, 1, len(response.OperationResults))
|
||||||
|
|
||||||
|
operationResponse := response.OperationResults[0]
|
||||||
|
assert.NotNil(t, operationResponse)
|
||||||
|
assert.NotNil(t, operationResponse.ResourceBody)
|
||||||
|
assert.Equal(t, int32(http.StatusCreated), operationResponse.StatusCode)
|
||||||
|
|
||||||
|
var itemResponseBody map[string]interface{}
|
||||||
|
json.Unmarshal(operationResponse.ResourceBody, &itemResponseBody)
|
||||||
|
assert.Equal(t, newItem["id"], itemResponseBody["id"])
|
||||||
|
assert.Equal(t, newItem["pk"], itemResponseBody["pk"])
|
||||||
|
|
||||||
|
updatedDoc, _ := ts.Repository.GetDocument(testDatabaseName, testCollectionName, newItem["id"].(string))
|
||||||
|
assert.Equal(t, newItem["id"], updatedDoc["id"])
|
||||||
|
assert.Equal(t, newItem["pk"], updatedDoc["pk"])
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("Should execute UPSERT transactional batch", func(t *testing.T) {
|
||||||
|
context := context.TODO()
|
||||||
|
batch := collectionClient.NewTransactionalBatch(azcosmos.NewPartitionKeyString("pk"))
|
||||||
|
|
||||||
|
newItem := map[string]interface{}{
|
||||||
|
"id": "678901",
|
||||||
|
"pk": "666",
|
||||||
|
}
|
||||||
|
bytes, err := json.Marshal(newItem)
|
||||||
|
assert.Nil(t, err)
|
||||||
|
|
||||||
|
batch.UpsertItem(bytes, nil)
|
||||||
|
response, err := collectionClient.ExecuteTransactionalBatch(context, batch, &azcosmos.TransactionalBatchOptions{})
|
||||||
|
assert.Nil(t, err)
|
||||||
|
assert.True(t, response.Success)
|
||||||
|
assert.Equal(t, 1, len(response.OperationResults))
|
||||||
|
|
||||||
|
operationResponse := response.OperationResults[0]
|
||||||
|
assert.NotNil(t, operationResponse)
|
||||||
|
assert.NotNil(t, operationResponse.ResourceBody)
|
||||||
|
assert.Equal(t, int32(http.StatusCreated), operationResponse.StatusCode)
|
||||||
|
|
||||||
|
var itemResponseBody map[string]interface{}
|
||||||
|
json.Unmarshal(operationResponse.ResourceBody, &itemResponseBody)
|
||||||
|
assert.Equal(t, newItem["id"], itemResponseBody["id"])
|
||||||
|
assert.Equal(t, newItem["pk"], itemResponseBody["pk"])
|
||||||
|
|
||||||
|
updatedDoc, _ := ts.Repository.GetDocument(testDatabaseName, testCollectionName, newItem["id"].(string))
|
||||||
|
assert.Equal(t, newItem["id"], updatedDoc["id"])
|
||||||
|
assert.Equal(t, newItem["pk"], updatedDoc["pk"])
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("Should execute READ transactional batch", func(t *testing.T) {
|
||||||
|
context := context.TODO()
|
||||||
|
batch := collectionClient.NewTransactionalBatch(azcosmos.NewPartitionKeyString("pk"))
|
||||||
|
|
||||||
|
batch.ReadItem("67890", nil)
|
||||||
|
response, err := collectionClient.ExecuteTransactionalBatch(context, batch, &azcosmos.TransactionalBatchOptions{})
|
||||||
|
assert.Nil(t, err)
|
||||||
|
assert.True(t, response.Success)
|
||||||
|
assert.Equal(t, 1, len(response.OperationResults))
|
||||||
|
|
||||||
|
operationResponse := response.OperationResults[0]
|
||||||
|
assert.NotNil(t, operationResponse)
|
||||||
|
assert.NotNil(t, operationResponse.ResourceBody)
|
||||||
|
assert.Equal(t, int32(http.StatusOK), operationResponse.StatusCode)
|
||||||
|
|
||||||
|
var itemResponseBody map[string]interface{}
|
||||||
|
json.Unmarshal(operationResponse.ResourceBody, &itemResponseBody)
|
||||||
|
assert.Equal(t, "67890", itemResponseBody["id"])
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -204,14 +204,18 @@ Cosmium strives to support the core features of Cosmos DB, including:
|
|||||||
| IS_PRIMITIVE | Yes |
|
| IS_PRIMITIVE | Yes |
|
||||||
| IS_STRING | Yes |
|
| IS_STRING | Yes |
|
||||||
|
|
||||||
### Document Batch Requests
|
### Transactional batch operations
|
||||||
|
|
||||||
|
Note: There's actually no transaction here. Think of this as a 'bulk operation' that can partially succeed.
|
||||||
|
|
||||||
| Operation | Implemented |
|
| Operation | Implemented |
|
||||||
| --------- | ----------- |
|
| --------- | ----------- |
|
||||||
| Create | No |
|
| Create | Yes |
|
||||||
| Update | No |
|
| Delete | Yes |
|
||||||
| Delete | No |
|
| Replace | Yes |
|
||||||
| Read | No |
|
| Upsert | Yes |
|
||||||
|
| Read | Yes |
|
||||||
|
| Patch | No |
|
||||||
|
|
||||||
## Known Differences
|
## Known Differences
|
||||||
|
|
||||||
|
|||||||
22
go.mod
22
go.mod
@@ -1,6 +1,6 @@
|
|||||||
module github.com/pikami/cosmium
|
module github.com/pikami/cosmium
|
||||||
|
|
||||||
go 1.22.0
|
go 1.24.0
|
||||||
|
|
||||||
require (
|
require (
|
||||||
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.12.0
|
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.12.0
|
||||||
@@ -9,13 +9,13 @@ require (
|
|||||||
github.com/gin-gonic/gin v1.10.0
|
github.com/gin-gonic/gin v1.10.0
|
||||||
github.com/google/uuid v1.6.0
|
github.com/google/uuid v1.6.0
|
||||||
github.com/stretchr/testify v1.10.0
|
github.com/stretchr/testify v1.10.0
|
||||||
golang.org/x/exp v0.0.0-20250106191152-7588d65b2ba8
|
golang.org/x/exp v0.0.0-20250218142911-aa4b98e5adaa
|
||||||
)
|
)
|
||||||
|
|
||||||
require (
|
require (
|
||||||
github.com/Azure/azure-sdk-for-go v68.0.0+incompatible // indirect
|
github.com/Azure/azure-sdk-for-go v68.0.0+incompatible // indirect
|
||||||
github.com/Azure/azure-sdk-for-go/sdk/internal v1.10.0 // indirect
|
github.com/Azure/azure-sdk-for-go/sdk/internal v1.10.0 // indirect
|
||||||
github.com/bytedance/sonic v1.12.7 // indirect
|
github.com/bytedance/sonic v1.12.8 // indirect
|
||||||
github.com/bytedance/sonic/loader v0.2.3 // indirect
|
github.com/bytedance/sonic/loader v0.2.3 // indirect
|
||||||
github.com/cloudwego/base64x v0.1.5 // indirect
|
github.com/cloudwego/base64x v0.1.5 // indirect
|
||||||
github.com/davecgh/go-spew v1.1.1 // indirect
|
github.com/davecgh/go-spew v1.1.1 // indirect
|
||||||
@@ -23,8 +23,8 @@ require (
|
|||||||
github.com/gin-contrib/sse v1.0.0 // indirect
|
github.com/gin-contrib/sse v1.0.0 // indirect
|
||||||
github.com/go-playground/locales v0.14.1 // indirect
|
github.com/go-playground/locales v0.14.1 // indirect
|
||||||
github.com/go-playground/universal-translator v0.18.1 // indirect
|
github.com/go-playground/universal-translator v0.18.1 // indirect
|
||||||
github.com/go-playground/validator/v10 v10.24.0 // indirect
|
github.com/go-playground/validator/v10 v10.25.0 // indirect
|
||||||
github.com/goccy/go-json v0.10.4 // indirect
|
github.com/goccy/go-json v0.10.5 // indirect
|
||||||
github.com/json-iterator/go v1.1.12 // indirect
|
github.com/json-iterator/go v1.1.12 // indirect
|
||||||
github.com/klauspost/cpuid/v2 v2.2.9 // indirect
|
github.com/klauspost/cpuid/v2 v2.2.9 // indirect
|
||||||
github.com/leodido/go-urn v1.4.0 // indirect
|
github.com/leodido/go-urn v1.4.0 // indirect
|
||||||
@@ -36,11 +36,11 @@ require (
|
|||||||
github.com/pmezard/go-difflib v1.0.0 // indirect
|
github.com/pmezard/go-difflib v1.0.0 // indirect
|
||||||
github.com/twitchyliquid64/golang-asm v0.15.1 // indirect
|
github.com/twitchyliquid64/golang-asm v0.15.1 // indirect
|
||||||
github.com/ugorji/go/codec v1.2.12 // indirect
|
github.com/ugorji/go/codec v1.2.12 // indirect
|
||||||
golang.org/x/arch v0.13.0 // indirect
|
golang.org/x/arch v0.14.0 // indirect
|
||||||
golang.org/x/crypto v0.32.0 // indirect
|
golang.org/x/crypto v0.33.0 // indirect
|
||||||
golang.org/x/net v0.34.0 // indirect
|
golang.org/x/net v0.35.0 // indirect
|
||||||
golang.org/x/sys v0.29.0 // indirect
|
golang.org/x/sys v0.30.0 // indirect
|
||||||
golang.org/x/text v0.21.0 // indirect
|
golang.org/x/text v0.22.0 // indirect
|
||||||
google.golang.org/protobuf v1.36.4 // indirect
|
google.golang.org/protobuf v1.36.5 // indirect
|
||||||
gopkg.in/yaml.v3 v3.0.1 // indirect
|
gopkg.in/yaml.v3 v3.0.1 // indirect
|
||||||
)
|
)
|
||||||
|
|||||||
40
go.sum
40
go.sum
@@ -10,8 +10,8 @@ github.com/Azure/azure-sdk-for-go/sdk/internal v1.10.0 h1:ywEEhmNahHBihViHepv3xP
|
|||||||
github.com/Azure/azure-sdk-for-go/sdk/internal v1.10.0/go.mod h1:iZDifYGJTIgIIkYRNWPENUnqx6bJ2xnSDFI2tjwZNuY=
|
github.com/Azure/azure-sdk-for-go/sdk/internal v1.10.0/go.mod h1:iZDifYGJTIgIIkYRNWPENUnqx6bJ2xnSDFI2tjwZNuY=
|
||||||
github.com/AzureAD/microsoft-authentication-library-for-go v1.2.2 h1:XHOnouVk1mxXfQidrMEnLlPk9UMeRtyBTnEFtxkV0kU=
|
github.com/AzureAD/microsoft-authentication-library-for-go v1.2.2 h1:XHOnouVk1mxXfQidrMEnLlPk9UMeRtyBTnEFtxkV0kU=
|
||||||
github.com/AzureAD/microsoft-authentication-library-for-go v1.2.2/go.mod h1:wP83P5OoQ5p6ip3ScPr0BAq0BvuPAvacpEuSzyouqAI=
|
github.com/AzureAD/microsoft-authentication-library-for-go v1.2.2/go.mod h1:wP83P5OoQ5p6ip3ScPr0BAq0BvuPAvacpEuSzyouqAI=
|
||||||
github.com/bytedance/sonic v1.12.7 h1:CQU8pxOy9HToxhndH0Kx/S1qU/CuS9GnKYrGioDcU1Q=
|
github.com/bytedance/sonic v1.12.8 h1:4xYRVRlXIgvSZ4e8iVTlMF5szgpXd4AfvuWgA8I8lgs=
|
||||||
github.com/bytedance/sonic v1.12.7/go.mod h1:tnbal4mxOMju17EGfknm2XyYcpyCnIROYOEYuemj13I=
|
github.com/bytedance/sonic v1.12.8/go.mod h1:uVvFidNmlt9+wa31S1urfwwthTWteBgG0hWuoKAXTx8=
|
||||||
github.com/bytedance/sonic/loader v0.1.1/go.mod h1:ncP89zfokxS5LZrJxl5z0UJcsk4M4yY2JpfqGeCtNLU=
|
github.com/bytedance/sonic/loader v0.1.1/go.mod h1:ncP89zfokxS5LZrJxl5z0UJcsk4M4yY2JpfqGeCtNLU=
|
||||||
github.com/bytedance/sonic/loader v0.2.3 h1:yctD0Q3v2NOGfSWPLPvG2ggA2kV6TS6s4wioyEqssH0=
|
github.com/bytedance/sonic/loader v0.2.3 h1:yctD0Q3v2NOGfSWPLPvG2ggA2kV6TS6s4wioyEqssH0=
|
||||||
github.com/bytedance/sonic/loader v0.2.3/go.mod h1:N8A3vUdtUebEY2/VQC0MyhYeKUFosQU6FxH2JmUe6VI=
|
github.com/bytedance/sonic/loader v0.2.3/go.mod h1:N8A3vUdtUebEY2/VQC0MyhYeKUFosQU6FxH2JmUe6VI=
|
||||||
@@ -35,10 +35,10 @@ github.com/go-playground/locales v0.14.1 h1:EWaQ/wswjilfKLTECiXz7Rh+3BjFhfDFKv/o
|
|||||||
github.com/go-playground/locales v0.14.1/go.mod h1:hxrqLVvrK65+Rwrd5Fc6F2O76J/NuW9t0sjnWqG1slY=
|
github.com/go-playground/locales v0.14.1/go.mod h1:hxrqLVvrK65+Rwrd5Fc6F2O76J/NuW9t0sjnWqG1slY=
|
||||||
github.com/go-playground/universal-translator v0.18.1 h1:Bcnm0ZwsGyWbCzImXv+pAJnYK9S473LQFuzCbDbfSFY=
|
github.com/go-playground/universal-translator v0.18.1 h1:Bcnm0ZwsGyWbCzImXv+pAJnYK9S473LQFuzCbDbfSFY=
|
||||||
github.com/go-playground/universal-translator v0.18.1/go.mod h1:xekY+UJKNuX9WP91TpwSH2VMlDf28Uj24BCp08ZFTUY=
|
github.com/go-playground/universal-translator v0.18.1/go.mod h1:xekY+UJKNuX9WP91TpwSH2VMlDf28Uj24BCp08ZFTUY=
|
||||||
github.com/go-playground/validator/v10 v10.24.0 h1:KHQckvo8G6hlWnrPX4NJJ+aBfWNAE/HH+qdL2cBpCmg=
|
github.com/go-playground/validator/v10 v10.25.0 h1:5Dh7cjvzR7BRZadnsVOzPhWsrwUr0nmsZJxEAnFLNO8=
|
||||||
github.com/go-playground/validator/v10 v10.24.0/go.mod h1:GGzBIJMuE98Ic/kJsBXbz1x/7cByt++cQ+YOuDM5wus=
|
github.com/go-playground/validator/v10 v10.25.0/go.mod h1:GGzBIJMuE98Ic/kJsBXbz1x/7cByt++cQ+YOuDM5wus=
|
||||||
github.com/goccy/go-json v0.10.4 h1:JSwxQzIqKfmFX1swYPpUThQZp/Ka4wzJdK0LWVytLPM=
|
github.com/goccy/go-json v0.10.5 h1:Fq85nIqj+gXn/S5ahsiTlK3TmC85qgirsdTP/+DeaC4=
|
||||||
github.com/goccy/go-json v0.10.4/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M=
|
github.com/goccy/go-json v0.10.5/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M=
|
||||||
github.com/golang-jwt/jwt v3.2.1+incompatible h1:73Z+4BJcrTC+KczS6WvTPvRGOp1WmfEP4Q1lOd9Z/+c=
|
github.com/golang-jwt/jwt v3.2.1+incompatible h1:73Z+4BJcrTC+KczS6WvTPvRGOp1WmfEP4Q1lOd9Z/+c=
|
||||||
github.com/golang-jwt/jwt/v5 v5.2.1 h1:OuVbFODueb089Lh128TAcimifWaLhJwVflnrgM17wHk=
|
github.com/golang-jwt/jwt/v5 v5.2.1 h1:OuVbFODueb089Lh128TAcimifWaLhJwVflnrgM17wHk=
|
||||||
github.com/golang-jwt/jwt/v5 v5.2.1/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk=
|
github.com/golang-jwt/jwt/v5 v5.2.1/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk=
|
||||||
@@ -94,21 +94,21 @@ github.com/twitchyliquid64/golang-asm v0.15.1 h1:SU5vSMR7hnwNxj24w34ZyCi/FmDZTkS
|
|||||||
github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2+aY1QWCk3Cedj/Gdt08=
|
github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2+aY1QWCk3Cedj/Gdt08=
|
||||||
github.com/ugorji/go/codec v1.2.12 h1:9LC83zGrHhuUA9l16C9AHXAqEV/2wBQ4nkvumAE65EE=
|
github.com/ugorji/go/codec v1.2.12 h1:9LC83zGrHhuUA9l16C9AHXAqEV/2wBQ4nkvumAE65EE=
|
||||||
github.com/ugorji/go/codec v1.2.12/go.mod h1:UNopzCgEMSXjBc6AOMqYvWC1ktqTAfzJZUZgYf6w6lg=
|
github.com/ugorji/go/codec v1.2.12/go.mod h1:UNopzCgEMSXjBc6AOMqYvWC1ktqTAfzJZUZgYf6w6lg=
|
||||||
golang.org/x/arch v0.13.0 h1:KCkqVVV1kGg0X87TFysjCJ8MxtZEIU4Ja/yXGeoECdA=
|
golang.org/x/arch v0.14.0 h1:z9JUEZWr8x4rR0OU6c4/4t6E6jOZ8/QBS2bBYBm4tx4=
|
||||||
golang.org/x/arch v0.13.0/go.mod h1:FEVrYAQjsQXMVJ1nsMoVVXPZg6p2JE2mx8psSWTDQys=
|
golang.org/x/arch v0.14.0/go.mod h1:FEVrYAQjsQXMVJ1nsMoVVXPZg6p2JE2mx8psSWTDQys=
|
||||||
golang.org/x/crypto v0.32.0 h1:euUpcYgM8WcP71gNpTqQCn6rC2t6ULUPiOzfWaXVVfc=
|
golang.org/x/crypto v0.33.0 h1:IOBPskki6Lysi0lo9qQvbxiQ+FvsCC/YWOecCHAixus=
|
||||||
golang.org/x/crypto v0.32.0/go.mod h1:ZnnJkOaASj8g0AjIduWNlq2NRxL0PlBrbKVyZ6V/Ugc=
|
golang.org/x/crypto v0.33.0/go.mod h1:bVdXmD7IV/4GdElGPozy6U7lWdRXA4qyRVGJV57uQ5M=
|
||||||
golang.org/x/exp v0.0.0-20250106191152-7588d65b2ba8 h1:yqrTHse8TCMW1M1ZCP+VAR/l0kKxwaAIqN/il7x4voA=
|
golang.org/x/exp v0.0.0-20250218142911-aa4b98e5adaa h1:t2QcU6V556bFjYgu4L6C+6VrCPyJZ+eyRsABUPs1mz4=
|
||||||
golang.org/x/exp v0.0.0-20250106191152-7588d65b2ba8/go.mod h1:tujkw807nyEEAamNbDrEGzRav+ilXA7PCRAd6xsmwiU=
|
golang.org/x/exp v0.0.0-20250218142911-aa4b98e5adaa/go.mod h1:BHOTPb3L19zxehTsLoJXVaTktb06DFgmdW6Wb9s8jqk=
|
||||||
golang.org/x/net v0.34.0 h1:Mb7Mrk043xzHgnRM88suvJFwzVrRfHEHJEl5/71CKw0=
|
golang.org/x/net v0.35.0 h1:T5GQRQb2y08kTAByq9L4/bz8cipCdA8FbRTXewonqY8=
|
||||||
golang.org/x/net v0.34.0/go.mod h1:di0qlW3YNM5oh6GqDGQr92MyTozJPmybPK4Ev/Gm31k=
|
golang.org/x/net v0.35.0/go.mod h1:EglIi67kWsHKlRzzVMUD93VMSWGFOMSZgxFjparz1Qk=
|
||||||
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||||
golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU=
|
golang.org/x/sys v0.30.0 h1:QjkSwP/36a20jFYWkSue1YwXzLmsV5Gfq7Eiy72C1uc=
|
||||||
golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||||
golang.org/x/text v0.21.0 h1:zyQAAkrwaneQ066sspRyJaG9VNi/YJ1NfzcGB3hZ/qo=
|
golang.org/x/text v0.22.0 h1:bofq7m3/HAFvbF51jz3Q9wLg3jkvSPuiZu/pD1XwgtM=
|
||||||
golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ=
|
golang.org/x/text v0.22.0/go.mod h1:YRoo4H8PVmsu+E3Ou7cqLVH8oXWIHVoX0jqUWALQhfY=
|
||||||
google.golang.org/protobuf v1.36.4 h1:6A3ZDJHn/eNqc1i+IdefRzy/9PokBTPvcqMySR7NNIM=
|
google.golang.org/protobuf v1.36.5 h1:tPhr+woSbjfYvY6/GPufUoYizxw1cF/yFoxJ2fmpwlM=
|
||||||
google.golang.org/protobuf v1.36.4/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE=
|
google.golang.org/protobuf v1.36.5/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE=
|
||||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||||
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
|
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
|
||||||
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
|
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
|
||||||
|
|||||||
@@ -50,6 +50,10 @@ func (r *DataRepository) DeleteCollection(databaseId string, collectionId string
|
|||||||
}
|
}
|
||||||
|
|
||||||
delete(r.storeState.Collections[databaseId], collectionId)
|
delete(r.storeState.Collections[databaseId], collectionId)
|
||||||
|
delete(r.storeState.Documents[databaseId], collectionId)
|
||||||
|
delete(r.storeState.Triggers[databaseId], collectionId)
|
||||||
|
delete(r.storeState.StoredProcedures[databaseId], collectionId)
|
||||||
|
delete(r.storeState.UserDefinedFunctions[databaseId], collectionId)
|
||||||
|
|
||||||
return repositorymodels.StatusOk
|
return repositorymodels.StatusOk
|
||||||
}
|
}
|
||||||
@@ -71,7 +75,7 @@ func (r *DataRepository) CreateCollection(databaseId string, newCollection repos
|
|||||||
newCollection = structhidrators.Hidrate(newCollection).(repositorymodels.Collection)
|
newCollection = structhidrators.Hidrate(newCollection).(repositorymodels.Collection)
|
||||||
|
|
||||||
newCollection.TimeStamp = time.Now().Unix()
|
newCollection.TimeStamp = time.Now().Unix()
|
||||||
newCollection.ResourceID = resourceid.NewCombined(database.ResourceID, resourceid.New())
|
newCollection.ResourceID = resourceid.NewCombined(database.ResourceID, resourceid.New(resourceid.ResourceTypeCollection))
|
||||||
newCollection.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
newCollection.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
||||||
newCollection.Self = fmt.Sprintf("dbs/%s/colls/%s/", database.ResourceID, newCollection.ResourceID)
|
newCollection.Self = fmt.Sprintf("dbs/%s/colls/%s/", database.ResourceID, newCollection.ResourceID)
|
||||||
|
|
||||||
|
|||||||
@@ -37,6 +37,11 @@ func (r *DataRepository) DeleteDatabase(id string) repositorymodels.RepositorySt
|
|||||||
}
|
}
|
||||||
|
|
||||||
delete(r.storeState.Databases, id)
|
delete(r.storeState.Databases, id)
|
||||||
|
delete(r.storeState.Collections, id)
|
||||||
|
delete(r.storeState.Documents, id)
|
||||||
|
delete(r.storeState.Triggers, id)
|
||||||
|
delete(r.storeState.StoredProcedures, id)
|
||||||
|
delete(r.storeState.UserDefinedFunctions, id)
|
||||||
|
|
||||||
return repositorymodels.StatusOk
|
return repositorymodels.StatusOk
|
||||||
}
|
}
|
||||||
@@ -50,7 +55,7 @@ func (r *DataRepository) CreateDatabase(newDatabase repositorymodels.Database) (
|
|||||||
}
|
}
|
||||||
|
|
||||||
newDatabase.TimeStamp = time.Now().Unix()
|
newDatabase.TimeStamp = time.Now().Unix()
|
||||||
newDatabase.ResourceID = resourceid.New()
|
newDatabase.ResourceID = resourceid.New(resourceid.ResourceTypeDatabase)
|
||||||
newDatabase.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
newDatabase.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
||||||
newDatabase.Self = fmt.Sprintf("dbs/%s/", newDatabase.ResourceID)
|
newDatabase.Self = fmt.Sprintf("dbs/%s/", newDatabase.ResourceID)
|
||||||
|
|
||||||
|
|||||||
@@ -95,7 +95,7 @@ func (r *DataRepository) CreateDocument(databaseId string, collectionId string,
|
|||||||
}
|
}
|
||||||
|
|
||||||
document["_ts"] = time.Now().Unix()
|
document["_ts"] = time.Now().Unix()
|
||||||
document["_rid"] = resourceid.NewCombined(database.ResourceID, collection.ResourceID, resourceid.New())
|
document["_rid"] = resourceid.NewCombined(collection.ResourceID, resourceid.New(resourceid.ResourceTypeDocument))
|
||||||
document["_etag"] = fmt.Sprintf("\"%s\"", uuid.New())
|
document["_etag"] = fmt.Sprintf("\"%s\"", uuid.New())
|
||||||
document["_self"] = fmt.Sprintf("dbs/%s/colls/%s/docs/%s/", database.ResourceID, collection.ResourceID, document["_rid"])
|
document["_self"] = fmt.Sprintf("dbs/%s/colls/%s/docs/%s/", database.ResourceID, collection.ResourceID, document["_rid"])
|
||||||
|
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ func (r *DataRepository) GetPartitionKeyRanges(databaseId string, collectionId s
|
|||||||
timestamp = collection.TimeStamp
|
timestamp = collection.TimeStamp
|
||||||
}
|
}
|
||||||
|
|
||||||
pkrResourceId := resourceid.NewCombined(databaseRid, collectionRid, resourceid.New())
|
pkrResourceId := resourceid.NewCombined(collectionRid, resourceid.New(resourceid.ResourceTypePartitionKeyRange))
|
||||||
pkrSelf := fmt.Sprintf("dbs/%s/colls/%s/pkranges/%s/", databaseRid, collectionRid, pkrResourceId)
|
pkrSelf := fmt.Sprintf("dbs/%s/colls/%s/pkranges/%s/", databaseRid, collectionRid, pkrResourceId)
|
||||||
etag := fmt.Sprintf("\"%s\"", uuid.New())
|
etag := fmt.Sprintf("\"%s\"", uuid.New())
|
||||||
|
|
||||||
|
|||||||
@@ -81,7 +81,7 @@ func (r *DataRepository) CreateStoredProcedure(databaseId string, collectionId s
|
|||||||
}
|
}
|
||||||
|
|
||||||
sp.TimeStamp = time.Now().Unix()
|
sp.TimeStamp = time.Now().Unix()
|
||||||
sp.ResourceID = resourceid.NewCombined(database.ResourceID, collection.ResourceID, resourceid.New())
|
sp.ResourceID = resourceid.NewCombined(collection.ResourceID, resourceid.New(resourceid.ResourceTypeStoredProcedure))
|
||||||
sp.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
sp.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
||||||
sp.Self = fmt.Sprintf("dbs/%s/colls/%s/sprocs/%s/", database.ResourceID, collection.ResourceID, sp.ResourceID)
|
sp.Self = fmt.Sprintf("dbs/%s/colls/%s/sprocs/%s/", database.ResourceID, collection.ResourceID, sp.ResourceID)
|
||||||
|
|
||||||
|
|||||||
@@ -81,7 +81,7 @@ func (r *DataRepository) CreateTrigger(databaseId string, collectionId string, t
|
|||||||
}
|
}
|
||||||
|
|
||||||
trigger.TimeStamp = time.Now().Unix()
|
trigger.TimeStamp = time.Now().Unix()
|
||||||
trigger.ResourceID = resourceid.NewCombined(database.ResourceID, collection.ResourceID, resourceid.New())
|
trigger.ResourceID = resourceid.NewCombined(collection.ResourceID, resourceid.New(resourceid.ResourceTypeTrigger))
|
||||||
trigger.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
trigger.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
||||||
trigger.Self = fmt.Sprintf("dbs/%s/colls/%s/triggers/%s/", database.ResourceID, collection.ResourceID, trigger.ResourceID)
|
trigger.Self = fmt.Sprintf("dbs/%s/colls/%s/triggers/%s/", database.ResourceID, collection.ResourceID, trigger.ResourceID)
|
||||||
|
|
||||||
|
|||||||
@@ -81,7 +81,7 @@ func (r *DataRepository) CreateUserDefinedFunction(databaseId string, collection
|
|||||||
}
|
}
|
||||||
|
|
||||||
udf.TimeStamp = time.Now().Unix()
|
udf.TimeStamp = time.Now().Unix()
|
||||||
udf.ResourceID = resourceid.NewCombined(database.ResourceID, collection.ResourceID, resourceid.New())
|
udf.ResourceID = resourceid.NewCombined(collection.ResourceID, resourceid.New(resourceid.ResourceTypeUserDefinedFunction))
|
||||||
udf.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
udf.ETag = fmt.Sprintf("\"%s\"", uuid.New())
|
||||||
udf.Self = fmt.Sprintf("dbs/%s/colls/%s/udfs/%s/", database.ResourceID, collection.ResourceID, udf.ResourceID)
|
udf.Self = fmt.Sprintf("dbs/%s/colls/%s/udfs/%s/", database.ResourceID, collection.ResourceID, udf.ResourceID)
|
||||||
|
|
||||||
|
|||||||
@@ -3,32 +3,76 @@ package resourceid
|
|||||||
import (
|
import (
|
||||||
"encoding/base64"
|
"encoding/base64"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
|
"strings"
|
||||||
|
|
||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
)
|
)
|
||||||
|
|
||||||
func New() string {
|
type ResourceType int
|
||||||
id := uuid.New().ID()
|
|
||||||
idBytes := uintToBytes(id)
|
|
||||||
|
|
||||||
// first byte should be bigger than 0x80 for collection ids
|
const (
|
||||||
// clients classify this id as "user" otherwise
|
ResourceTypeDatabase ResourceType = iota
|
||||||
if (idBytes[0] & 0x80) <= 0 {
|
ResourceTypeCollection
|
||||||
idBytes[0] = byte(rand.Intn(0x80) + 0x80)
|
ResourceTypeDocument
|
||||||
|
ResourceTypeStoredProcedure
|
||||||
|
ResourceTypeTrigger
|
||||||
|
ResourceTypeUserDefinedFunction
|
||||||
|
ResourceTypeConflict
|
||||||
|
ResourceTypePartitionKeyRange
|
||||||
|
ResourceTypeSchema
|
||||||
|
)
|
||||||
|
|
||||||
|
func New(resourceType ResourceType) string {
|
||||||
|
var idBytes []byte
|
||||||
|
switch resourceType {
|
||||||
|
case ResourceTypeDatabase:
|
||||||
|
idBytes = randomBytes(4)
|
||||||
|
case ResourceTypeCollection:
|
||||||
|
idBytes = randomBytes(4)
|
||||||
|
// first byte should be bigger than 0x80 for collection ids
|
||||||
|
// clients classify this id as "user" otherwise
|
||||||
|
if (idBytes[0] & 0x80) <= 0 {
|
||||||
|
idBytes[0] = byte(rand.Intn(0x80) + 0x80)
|
||||||
|
}
|
||||||
|
case ResourceTypeDocument:
|
||||||
|
idBytes = randomBytes(8)
|
||||||
|
idBytes[7] = byte(rand.Intn(0x10)) // Upper 4 bits = 0
|
||||||
|
case ResourceTypeStoredProcedure:
|
||||||
|
idBytes = randomBytes(8)
|
||||||
|
idBytes[7] = byte(rand.Intn(0x10)) | 0x08 // Upper 4 bits = 0x08
|
||||||
|
case ResourceTypeTrigger:
|
||||||
|
idBytes = randomBytes(8)
|
||||||
|
idBytes[7] = byte(rand.Intn(0x10)) | 0x07 // Upper 4 bits = 0x07
|
||||||
|
case ResourceTypeUserDefinedFunction:
|
||||||
|
idBytes = randomBytes(8)
|
||||||
|
idBytes[7] = byte(rand.Intn(0x10)) | 0x06 // Upper 4 bits = 0x06
|
||||||
|
case ResourceTypeConflict:
|
||||||
|
idBytes = randomBytes(8)
|
||||||
|
idBytes[7] = byte(rand.Intn(0x10)) | 0x04 // Upper 4 bits = 0x04
|
||||||
|
case ResourceTypePartitionKeyRange:
|
||||||
|
// we don't do partitions yet, so just use a fixed id
|
||||||
|
idBytes = []byte{0x69, 0x69, 0x69, 0x69, 0x69, 0x69, 0x69, 0x50}
|
||||||
|
case ResourceTypeSchema:
|
||||||
|
idBytes = randomBytes(8)
|
||||||
|
idBytes[7] = byte(rand.Intn(0x10)) | 0x09 // Upper 4 bits = 0x09
|
||||||
|
default:
|
||||||
|
idBytes = randomBytes(4)
|
||||||
}
|
}
|
||||||
|
|
||||||
return base64.StdEncoding.EncodeToString(idBytes)
|
encoded := base64.StdEncoding.EncodeToString(idBytes)
|
||||||
|
return strings.ReplaceAll(encoded, "/", "-")
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewCombined(ids ...string) string {
|
func NewCombined(ids ...string) string {
|
||||||
combinedIdBytes := make([]byte, 0)
|
combinedIdBytes := make([]byte, 0)
|
||||||
|
|
||||||
for _, id := range ids {
|
for _, id := range ids {
|
||||||
idBytes, _ := base64.StdEncoding.DecodeString(id)
|
idBytes, _ := base64.StdEncoding.DecodeString(strings.ReplaceAll(id, "-", "/"))
|
||||||
combinedIdBytes = append(combinedIdBytes, idBytes...)
|
combinedIdBytes = append(combinedIdBytes, idBytes...)
|
||||||
}
|
}
|
||||||
|
|
||||||
return base64.StdEncoding.EncodeToString(combinedIdBytes)
|
encoded := base64.StdEncoding.EncodeToString(combinedIdBytes)
|
||||||
|
return strings.ReplaceAll(encoded, "/", "-")
|
||||||
}
|
}
|
||||||
|
|
||||||
func uintToBytes(id uint32) []byte {
|
func uintToBytes(id uint32) []byte {
|
||||||
@@ -39,3 +83,13 @@ func uintToBytes(id uint32) []byte {
|
|||||||
|
|
||||||
return buf
|
return buf
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func randomBytes(count int) []byte {
|
||||||
|
buf := make([]byte, count)
|
||||||
|
for i := 0; i < count; i += 4 {
|
||||||
|
id := uuid.New().ID()
|
||||||
|
idBytes := uintToBytes(id)
|
||||||
|
copy(buf[i:], idBytes)
|
||||||
|
}
|
||||||
|
return buf
|
||||||
|
}
|
||||||
|
|||||||
@@ -2107,40 +2107,40 @@ var g = &grammar{
|
|||||||
alternatives: []any{
|
alternatives: []any{
|
||||||
&litMatcher{
|
&litMatcher{
|
||||||
pos: position{line: 406, col: 24, offset: 11503},
|
pos: position{line: 406, col: 24, offset: 11503},
|
||||||
val: "=",
|
|
||||||
ignoreCase: false,
|
|
||||||
want: "\"=\"",
|
|
||||||
},
|
|
||||||
&litMatcher{
|
|
||||||
pos: position{line: 406, col: 30, offset: 11509},
|
|
||||||
val: "!=",
|
|
||||||
ignoreCase: false,
|
|
||||||
want: "\"!=\"",
|
|
||||||
},
|
|
||||||
&litMatcher{
|
|
||||||
pos: position{line: 406, col: 37, offset: 11516},
|
|
||||||
val: "<",
|
|
||||||
ignoreCase: false,
|
|
||||||
want: "\"<\"",
|
|
||||||
},
|
|
||||||
&litMatcher{
|
|
||||||
pos: position{line: 406, col: 43, offset: 11522},
|
|
||||||
val: "<=",
|
val: "<=",
|
||||||
ignoreCase: false,
|
ignoreCase: false,
|
||||||
want: "\"<=\"",
|
want: "\"<=\"",
|
||||||
},
|
},
|
||||||
&litMatcher{
|
&litMatcher{
|
||||||
pos: position{line: 406, col: 50, offset: 11529},
|
pos: position{line: 406, col: 31, offset: 11510},
|
||||||
val: ">",
|
|
||||||
ignoreCase: false,
|
|
||||||
want: "\">\"",
|
|
||||||
},
|
|
||||||
&litMatcher{
|
|
||||||
pos: position{line: 406, col: 56, offset: 11535},
|
|
||||||
val: ">=",
|
val: ">=",
|
||||||
ignoreCase: false,
|
ignoreCase: false,
|
||||||
want: "\">=\"",
|
want: "\">=\"",
|
||||||
},
|
},
|
||||||
|
&litMatcher{
|
||||||
|
pos: position{line: 406, col: 38, offset: 11517},
|
||||||
|
val: "=",
|
||||||
|
ignoreCase: false,
|
||||||
|
want: "\"=\"",
|
||||||
|
},
|
||||||
|
&litMatcher{
|
||||||
|
pos: position{line: 406, col: 44, offset: 11523},
|
||||||
|
val: "!=",
|
||||||
|
ignoreCase: false,
|
||||||
|
want: "\"!=\"",
|
||||||
|
},
|
||||||
|
&litMatcher{
|
||||||
|
pos: position{line: 406, col: 51, offset: 11530},
|
||||||
|
val: "<",
|
||||||
|
ignoreCase: false,
|
||||||
|
want: "\"<\"",
|
||||||
|
},
|
||||||
|
&litMatcher{
|
||||||
|
pos: position{line: 406, col: 57, offset: 11536},
|
||||||
|
val: ">",
|
||||||
|
ignoreCase: false,
|
||||||
|
want: "\">\"",
|
||||||
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -403,7 +403,7 @@ OrderBy <- "ORDER"i ws "BY"i
|
|||||||
|
|
||||||
Offset <- "OFFSET"i
|
Offset <- "OFFSET"i
|
||||||
|
|
||||||
ComparisonOperator <- ("=" / "!=" / "<" / "<=" / ">" / ">=") {
|
ComparisonOperator <- ("<=" / ">=" / "=" / "!=" / "<" / ">") {
|
||||||
return string(c.text), nil
|
return string(c.text), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -67,7 +67,7 @@ func Test_Parse_Were(t *testing.T) {
|
|||||||
t,
|
t,
|
||||||
`select c.id
|
`select c.id
|
||||||
FROM c
|
FROM c
|
||||||
WHERE c.isCool=true AND (c.id = "123" OR c.id = "456")`,
|
WHERE c.isCool=true AND (c.id = "123" OR c.id <= "456")`,
|
||||||
parsers.SelectStmt{
|
parsers.SelectStmt{
|
||||||
SelectItems: []parsers.SelectItem{
|
SelectItems: []parsers.SelectItem{
|
||||||
{Path: []string{"c", "id"}},
|
{Path: []string{"c", "id"}},
|
||||||
@@ -90,7 +90,7 @@ func Test_Parse_Were(t *testing.T) {
|
|||||||
Right: testutils.SelectItem_Constant_String("123"),
|
Right: testutils.SelectItem_Constant_String("123"),
|
||||||
},
|
},
|
||||||
parsers.ComparisonExpression{
|
parsers.ComparisonExpression{
|
||||||
Operation: "=",
|
Operation: "<=",
|
||||||
Left: parsers.SelectItem{Path: []string{"c", "id"}},
|
Left: parsers.SelectItem{Path: []string{"c", "id"}},
|
||||||
Right: testutils.SelectItem_Constant_String("456"),
|
Right: testutils.SelectItem_Constant_String("456"),
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -60,6 +60,15 @@ func ExecuteQuery(query parsers.SelectStmt, documents []RowType) []RowType {
|
|||||||
projectedDocuments = deduplicate(projectedDocuments)
|
projectedDocuments = deduplicate(projectedDocuments)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Apply offset
|
||||||
|
if query.Offset > 0 {
|
||||||
|
if query.Offset < len(projectedDocuments) {
|
||||||
|
projectedDocuments = projectedDocuments[query.Offset:]
|
||||||
|
} else {
|
||||||
|
projectedDocuments = []RowType{}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Apply result limit
|
// Apply result limit
|
||||||
if query.Count > 0 && len(projectedDocuments) > query.Count {
|
if query.Count > 0 && len(projectedDocuments) > query.Count {
|
||||||
projectedDocuments = projectedDocuments[:query.Count]
|
projectedDocuments = projectedDocuments[:query.Count]
|
||||||
|
|||||||
@@ -10,10 +10,10 @@ import (
|
|||||||
|
|
||||||
func Test_Execute_Select(t *testing.T) {
|
func Test_Execute_Select(t *testing.T) {
|
||||||
mockData := []memoryexecutor.RowType{
|
mockData := []memoryexecutor.RowType{
|
||||||
map[string]interface{}{"id": "12345", "pk": 123, "_self": "self1", "_rid": "rid1", "_ts": 123456, "isCool": false},
|
map[string]interface{}{"id": "12345", "pk": 123, "_self": "self1", "_rid": "rid1", "_ts": 123456, "isCool": false, "order": 1},
|
||||||
map[string]interface{}{"id": "67890", "pk": 456, "_self": "self2", "_rid": "rid2", "_ts": 789012, "isCool": true},
|
map[string]interface{}{"id": "67890", "pk": 456, "_self": "self2", "_rid": "rid2", "_ts": 789012, "isCool": true, "order": 2},
|
||||||
map[string]interface{}{"id": "456", "pk": 456, "_self": "self2", "_rid": "rid2", "_ts": 789012, "isCool": true},
|
map[string]interface{}{"id": "456", "pk": 456, "_self": "self2", "_rid": "rid2", "_ts": 789012, "isCool": true, "order": 3},
|
||||||
map[string]interface{}{"id": "123", "pk": 456, "_self": "self2", "_rid": "rid2", "_ts": 789012, "isCool": true},
|
map[string]interface{}{"id": "123", "pk": 456, "_self": "self2", "_rid": "rid2", "_ts": 789012, "isCool": true, "order": 4},
|
||||||
}
|
}
|
||||||
|
|
||||||
t.Run("Should execute simple SELECT", func(t *testing.T) {
|
t.Run("Should execute simple SELECT", func(t *testing.T) {
|
||||||
@@ -108,15 +108,15 @@ func Test_Execute_Select(t *testing.T) {
|
|||||||
Offset: 1,
|
Offset: 1,
|
||||||
OrderExpressions: []parsers.OrderExpression{
|
OrderExpressions: []parsers.OrderExpression{
|
||||||
{
|
{
|
||||||
SelectItem: parsers.SelectItem{Path: []string{"c", "id"}},
|
SelectItem: parsers.SelectItem{Path: []string{"c", "order"}},
|
||||||
Direction: parsers.OrderDirectionDesc,
|
Direction: parsers.OrderDirectionDesc,
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
mockData,
|
mockData,
|
||||||
[]memoryexecutor.RowType{
|
[]memoryexecutor.RowType{
|
||||||
map[string]interface{}{"id": "67890", "pk": 456},
|
|
||||||
map[string]interface{}{"id": "456", "pk": 456},
|
map[string]interface{}{"id": "456", "pk": 456},
|
||||||
|
map[string]interface{}{"id": "67890", "pk": 456},
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -9,10 +9,14 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
func (r rowContext) strings_StringEquals(arguments []interface{}) bool {
|
func (r rowContext) strings_StringEquals(arguments []interface{}) bool {
|
||||||
str1 := r.parseString(arguments[0])
|
str1, str1ok := r.parseString(arguments[0])
|
||||||
str2 := r.parseString(arguments[1])
|
str2, str2ok := r.parseString(arguments[1])
|
||||||
ignoreCase := r.getBoolFlag(arguments)
|
ignoreCase := r.getBoolFlag(arguments)
|
||||||
|
|
||||||
|
if !str1ok || !str2ok {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
if ignoreCase {
|
if ignoreCase {
|
||||||
return strings.EqualFold(str1, str2)
|
return strings.EqualFold(str1, str2)
|
||||||
}
|
}
|
||||||
@@ -21,10 +25,14 @@ func (r rowContext) strings_StringEquals(arguments []interface{}) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_Contains(arguments []interface{}) bool {
|
func (r rowContext) strings_Contains(arguments []interface{}) bool {
|
||||||
str1 := r.parseString(arguments[0])
|
str1, str1ok := r.parseString(arguments[0])
|
||||||
str2 := r.parseString(arguments[1])
|
str2, str2ok := r.parseString(arguments[1])
|
||||||
ignoreCase := r.getBoolFlag(arguments)
|
ignoreCase := r.getBoolFlag(arguments)
|
||||||
|
|
||||||
|
if !str1ok || !str2ok {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
if ignoreCase {
|
if ignoreCase {
|
||||||
str1 = strings.ToLower(str1)
|
str1 = strings.ToLower(str1)
|
||||||
str2 = strings.ToLower(str2)
|
str2 = strings.ToLower(str2)
|
||||||
@@ -34,10 +42,14 @@ func (r rowContext) strings_Contains(arguments []interface{}) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_EndsWith(arguments []interface{}) bool {
|
func (r rowContext) strings_EndsWith(arguments []interface{}) bool {
|
||||||
str1 := r.parseString(arguments[0])
|
str1, str1ok := r.parseString(arguments[0])
|
||||||
str2 := r.parseString(arguments[1])
|
str2, str2ok := r.parseString(arguments[1])
|
||||||
ignoreCase := r.getBoolFlag(arguments)
|
ignoreCase := r.getBoolFlag(arguments)
|
||||||
|
|
||||||
|
if !str1ok || !str2ok {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
if ignoreCase {
|
if ignoreCase {
|
||||||
str1 = strings.ToLower(str1)
|
str1 = strings.ToLower(str1)
|
||||||
str2 = strings.ToLower(str2)
|
str2 = strings.ToLower(str2)
|
||||||
@@ -47,10 +59,14 @@ func (r rowContext) strings_EndsWith(arguments []interface{}) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_StartsWith(arguments []interface{}) bool {
|
func (r rowContext) strings_StartsWith(arguments []interface{}) bool {
|
||||||
str1 := r.parseString(arguments[0])
|
str1, str1ok := r.parseString(arguments[0])
|
||||||
str2 := r.parseString(arguments[1])
|
str2, str2ok := r.parseString(arguments[1])
|
||||||
ignoreCase := r.getBoolFlag(arguments)
|
ignoreCase := r.getBoolFlag(arguments)
|
||||||
|
|
||||||
|
if !str1ok || !str2ok {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
if ignoreCase {
|
if ignoreCase {
|
||||||
str1 = strings.ToLower(str1)
|
str1 = strings.ToLower(str1)
|
||||||
str2 = strings.ToLower(str2)
|
str2 = strings.ToLower(str2)
|
||||||
@@ -73,8 +89,12 @@ func (r rowContext) strings_Concat(arguments []interface{}) string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_IndexOf(arguments []interface{}) int {
|
func (r rowContext) strings_IndexOf(arguments []interface{}) int {
|
||||||
str1 := r.parseString(arguments[0])
|
str1, str1ok := r.parseString(arguments[0])
|
||||||
str2 := r.parseString(arguments[1])
|
str2, str2ok := r.parseString(arguments[1])
|
||||||
|
|
||||||
|
if !str1ok || !str2ok {
|
||||||
|
return -1
|
||||||
|
}
|
||||||
|
|
||||||
start := 0
|
start := 0
|
||||||
if len(arguments) > 2 && arguments[2] != nil {
|
if len(arguments) > 2 && arguments[2] != nil {
|
||||||
@@ -115,9 +135,13 @@ func (r rowContext) strings_Lower(arguments []interface{}) string {
|
|||||||
func (r rowContext) strings_Left(arguments []interface{}) string {
|
func (r rowContext) strings_Left(arguments []interface{}) string {
|
||||||
var ok bool
|
var ok bool
|
||||||
var length int
|
var length int
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
lengthEx := r.resolveSelectItem(arguments[1].(parsers.SelectItem))
|
lengthEx := r.resolveSelectItem(arguments[1].(parsers.SelectItem))
|
||||||
|
|
||||||
|
if !strOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
if length, ok = lengthEx.(int); !ok {
|
if length, ok = lengthEx.(int); !ok {
|
||||||
logger.ErrorLn("strings_Left - got parameters of wrong type")
|
logger.ErrorLn("strings_Left - got parameters of wrong type")
|
||||||
return ""
|
return ""
|
||||||
@@ -135,28 +159,45 @@ func (r rowContext) strings_Left(arguments []interface{}) string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_Length(arguments []interface{}) int {
|
func (r rowContext) strings_Length(arguments []interface{}) int {
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
|
if !strOk {
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
|
||||||
return len(str)
|
return len(str)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_LTrim(arguments []interface{}) string {
|
func (r rowContext) strings_LTrim(arguments []interface{}) string {
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
|
if !strOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
return strings.TrimLeft(str, " ")
|
return strings.TrimLeft(str, " ")
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_Replace(arguments []interface{}) string {
|
func (r rowContext) strings_Replace(arguments []interface{}) string {
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
oldStr := r.parseString(arguments[1])
|
oldStr, oldStrOk := r.parseString(arguments[1])
|
||||||
newStr := r.parseString(arguments[2])
|
newStr, newStrOk := r.parseString(arguments[2])
|
||||||
|
|
||||||
|
if !strOk || !oldStrOk || !newStrOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
return strings.Replace(str, oldStr, newStr, -1)
|
return strings.Replace(str, oldStr, newStr, -1)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_Replicate(arguments []interface{}) string {
|
func (r rowContext) strings_Replicate(arguments []interface{}) string {
|
||||||
var ok bool
|
var ok bool
|
||||||
var times int
|
var times int
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
timesEx := r.resolveSelectItem(arguments[1].(parsers.SelectItem))
|
timesEx := r.resolveSelectItem(arguments[1].(parsers.SelectItem))
|
||||||
|
|
||||||
|
if !strOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
if times, ok = timesEx.(int); !ok {
|
if times, ok = timesEx.(int); !ok {
|
||||||
logger.ErrorLn("strings_Replicate - got parameters of wrong type")
|
logger.ErrorLn("strings_Replicate - got parameters of wrong type")
|
||||||
return ""
|
return ""
|
||||||
@@ -174,9 +215,13 @@ func (r rowContext) strings_Replicate(arguments []interface{}) string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_Reverse(arguments []interface{}) string {
|
func (r rowContext) strings_Reverse(arguments []interface{}) string {
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
runes := []rune(str)
|
runes := []rune(str)
|
||||||
|
|
||||||
|
if !strOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
for i, j := 0, len(runes)-1; i < j; i, j = i+1, j-1 {
|
for i, j := 0, len(runes)-1; i < j; i, j = i+1, j-1 {
|
||||||
runes[i], runes[j] = runes[j], runes[i]
|
runes[i], runes[j] = runes[j], runes[i]
|
||||||
}
|
}
|
||||||
@@ -187,9 +232,13 @@ func (r rowContext) strings_Reverse(arguments []interface{}) string {
|
|||||||
func (r rowContext) strings_Right(arguments []interface{}) string {
|
func (r rowContext) strings_Right(arguments []interface{}) string {
|
||||||
var ok bool
|
var ok bool
|
||||||
var length int
|
var length int
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
lengthEx := r.resolveSelectItem(arguments[1].(parsers.SelectItem))
|
lengthEx := r.resolveSelectItem(arguments[1].(parsers.SelectItem))
|
||||||
|
|
||||||
|
if !strOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
if length, ok = lengthEx.(int); !ok {
|
if length, ok = lengthEx.(int); !ok {
|
||||||
logger.ErrorLn("strings_Right - got parameters of wrong type")
|
logger.ErrorLn("strings_Right - got parameters of wrong type")
|
||||||
return ""
|
return ""
|
||||||
@@ -207,7 +256,11 @@ func (r rowContext) strings_Right(arguments []interface{}) string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_RTrim(arguments []interface{}) string {
|
func (r rowContext) strings_RTrim(arguments []interface{}) string {
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
|
if !strOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
return strings.TrimRight(str, " ")
|
return strings.TrimRight(str, " ")
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -215,10 +268,14 @@ func (r rowContext) strings_Substring(arguments []interface{}) string {
|
|||||||
var ok bool
|
var ok bool
|
||||||
var startPos int
|
var startPos int
|
||||||
var length int
|
var length int
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
startPosEx := r.resolveSelectItem(arguments[1].(parsers.SelectItem))
|
startPosEx := r.resolveSelectItem(arguments[1].(parsers.SelectItem))
|
||||||
lengthEx := r.resolveSelectItem(arguments[2].(parsers.SelectItem))
|
lengthEx := r.resolveSelectItem(arguments[2].(parsers.SelectItem))
|
||||||
|
|
||||||
|
if !strOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
if startPos, ok = startPosEx.(int); !ok {
|
if startPos, ok = startPosEx.(int); !ok {
|
||||||
logger.ErrorLn("strings_Substring - got start parameters of wrong type")
|
logger.ErrorLn("strings_Substring - got start parameters of wrong type")
|
||||||
return ""
|
return ""
|
||||||
@@ -241,7 +298,11 @@ func (r rowContext) strings_Substring(arguments []interface{}) string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) strings_Trim(arguments []interface{}) string {
|
func (r rowContext) strings_Trim(arguments []interface{}) string {
|
||||||
str := r.parseString(arguments[0])
|
str, strOk := r.parseString(arguments[0])
|
||||||
|
if !strOk {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
return strings.TrimSpace(str)
|
return strings.TrimSpace(str)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -257,15 +318,15 @@ func (r rowContext) getBoolFlag(arguments []interface{}) bool {
|
|||||||
return ignoreCase
|
return ignoreCase
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r rowContext) parseString(argument interface{}) string {
|
func (r rowContext) parseString(argument interface{}) (value string, ok bool) {
|
||||||
exItem := argument.(parsers.SelectItem)
|
exItem := argument.(parsers.SelectItem)
|
||||||
ex := r.resolveSelectItem(exItem)
|
ex := r.resolveSelectItem(exItem)
|
||||||
if str1, ok := ex.(string); ok {
|
if str1, ok := ex.(string); ok {
|
||||||
return str1
|
return str1, true
|
||||||
}
|
}
|
||||||
|
|
||||||
logger.ErrorLn("StringEquals got parameters of wrong type")
|
logger.ErrorLn("StringEquals got parameters of wrong type")
|
||||||
return ""
|
return "", false
|
||||||
}
|
}
|
||||||
|
|
||||||
func convertToString(value interface{}) string {
|
func convertToString(value interface{}) string {
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package main
|
|||||||
import "C"
|
import "C"
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"strings"
|
||||||
|
|
||||||
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
||||||
)
|
)
|
||||||
@@ -20,7 +21,7 @@ func CreateCollection(serverName *C.char, databaseId *C.char, collectionJson *C.
|
|||||||
}
|
}
|
||||||
|
|
||||||
var collection repositorymodels.Collection
|
var collection repositorymodels.Collection
|
||||||
err := json.Unmarshal([]byte(collectionStr), &collection)
|
err := json.NewDecoder(strings.NewReader(collectionStr)).Decode(&collection)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return ResponseFailedToParseRequest
|
return ResponseFailedToParseRequest
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package main
|
|||||||
import "C"
|
import "C"
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"strings"
|
||||||
|
|
||||||
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
||||||
)
|
)
|
||||||
@@ -19,7 +20,7 @@ func CreateDatabase(serverName *C.char, databaseJson *C.char) int {
|
|||||||
}
|
}
|
||||||
|
|
||||||
var database repositorymodels.Database
|
var database repositorymodels.Database
|
||||||
err := json.Unmarshal([]byte(databaseStr), &database)
|
err := json.NewDecoder(strings.NewReader(databaseStr)).Decode(&database)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return ResponseFailedToParseRequest
|
return ResponseFailedToParseRequest
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package main
|
|||||||
import "C"
|
import "C"
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"strings"
|
||||||
|
|
||||||
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
repositorymodels "github.com/pikami/cosmium/internal/repository_models"
|
||||||
)
|
)
|
||||||
@@ -21,7 +22,7 @@ func CreateDocument(serverName *C.char, databaseId *C.char, collectionId *C.char
|
|||||||
}
|
}
|
||||||
|
|
||||||
var document repositorymodels.Document
|
var document repositorymodels.Document
|
||||||
err := json.Unmarshal([]byte(documentStr), &document)
|
err := json.NewDecoder(strings.NewReader(documentStr)).Decode(&document)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return ResponseFailedToParseRequest
|
return ResponseFailedToParseRequest
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,8 +13,10 @@ type ServerInstance struct {
|
|||||||
repository *repositories.DataRepository
|
repository *repositories.DataRepository
|
||||||
}
|
}
|
||||||
|
|
||||||
var serverInstances map[string]*ServerInstance
|
var (
|
||||||
var mutex sync.Mutex
|
serverInstances = make(map[string]*ServerInstance)
|
||||||
|
mutex = sync.Mutex{}
|
||||||
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
ResponseSuccess = 0
|
ResponseSuccess = 0
|
||||||
@@ -36,10 +38,6 @@ func getInstance(serverName string) (*ServerInstance, bool) {
|
|||||||
mutex.Lock()
|
mutex.Lock()
|
||||||
defer mutex.Unlock()
|
defer mutex.Unlock()
|
||||||
|
|
||||||
if serverInstances == nil {
|
|
||||||
serverInstances = make(map[string]*ServerInstance)
|
|
||||||
}
|
|
||||||
|
|
||||||
var ok bool
|
var ok bool
|
||||||
var serverInstance *ServerInstance
|
var serverInstance *ServerInstance
|
||||||
if serverInstance, ok = serverInstances[serverName]; !ok {
|
if serverInstance, ok = serverInstances[serverName]; !ok {
|
||||||
@@ -53,10 +51,6 @@ func addInstance(serverName string, serverInstance *ServerInstance) {
|
|||||||
mutex.Lock()
|
mutex.Lock()
|
||||||
defer mutex.Unlock()
|
defer mutex.Unlock()
|
||||||
|
|
||||||
if serverInstances == nil {
|
|
||||||
serverInstances = make(map[string]*ServerInstance)
|
|
||||||
}
|
|
||||||
|
|
||||||
serverInstances[serverName] = serverInstance
|
serverInstances[serverName] = serverInstance
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -64,10 +58,6 @@ func removeInstance(serverName string) {
|
|||||||
mutex.Lock()
|
mutex.Lock()
|
||||||
defer mutex.Unlock()
|
defer mutex.Unlock()
|
||||||
|
|
||||||
if serverInstances == nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
delete(serverInstances, serverName)
|
delete(serverInstances, serverName)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ package main
|
|||||||
import "C"
|
import "C"
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"strings"
|
||||||
"unsafe"
|
"unsafe"
|
||||||
|
|
||||||
"github.com/pikami/cosmium/api"
|
"github.com/pikami/cosmium/api"
|
||||||
@@ -15,21 +16,21 @@ import (
|
|||||||
|
|
||||||
//export CreateServerInstance
|
//export CreateServerInstance
|
||||||
func CreateServerInstance(serverName *C.char, configurationJSON *C.char) int {
|
func CreateServerInstance(serverName *C.char, configurationJSON *C.char) int {
|
||||||
configStr := C.GoString(configurationJSON)
|
|
||||||
serverNameStr := C.GoString(serverName)
|
serverNameStr := C.GoString(serverName)
|
||||||
|
configStr := C.GoString(configurationJSON)
|
||||||
|
|
||||||
if _, ok := getInstance(serverNameStr); ok {
|
if _, ok := getInstance(serverNameStr); ok {
|
||||||
return ResponseServerInstanceAlreadyExists
|
return ResponseServerInstanceAlreadyExists
|
||||||
}
|
}
|
||||||
|
|
||||||
var configuration config.ServerConfig
|
var configuration config.ServerConfig
|
||||||
err := json.Unmarshal([]byte(configStr), &configuration)
|
err := json.NewDecoder(strings.NewReader(configStr)).Decode(&configuration)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return ResponseFailedToParseConfiguration
|
return ResponseFailedToParseConfiguration
|
||||||
}
|
}
|
||||||
|
|
||||||
configuration.PopulateCalculatedFields()
|
|
||||||
configuration.ApplyDefaultsToEmptyFields()
|
configuration.ApplyDefaultsToEmptyFields()
|
||||||
|
configuration.PopulateCalculatedFields()
|
||||||
|
|
||||||
repository := repositories.NewDataRepository(repositories.RepositoryOptions{
|
repository := repositories.NewDataRepository(repositories.RepositoryOptions{
|
||||||
InitialDataFilePath: configuration.InitialDataFilePath,
|
InitialDataFilePath: configuration.InitialDataFilePath,
|
||||||
|
|||||||
Reference in New Issue
Block a user