Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pkg/auth/request.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ type ReadRequest struct {
type WriteRequest struct {
Principal Principal
Tenant TenantContext
Operation string // create | update
Operation string // create | update | patch
ResourceType string
ID string
RequiredPermissions []string
Expand Down
7 changes: 5 additions & 2 deletions pkg/core/batch.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"fmt"
"strings"

"github.com/degoke/health-ai-stack/pkg/hooks"
"github.com/degoke/health-ai-stack/pkg/types"
)

Expand Down Expand Up @@ -39,11 +40,13 @@ func (s *ResourceService) ProcessBatchBundle(ctx context.Context, bundle *types.
if err != nil {
return nil, exceptionErr("build batch response bundle", err)
}
return &types.ResourceEnvelope{
response := &types.ResourceEnvelope{
ResourceType: "Bundle",
JSON: responseJSON,
Hash: mustHash(responseJSON),
}, nil
}
s.runPostCommit(ctx, hooks.ActionBatch, response, nil)
return response, nil
}

type batchBundle struct {
Expand Down
39 changes: 19 additions & 20 deletions pkg/core/bundle.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"strings"
"time"

"github.com/degoke/health-ai-stack/pkg/hooks"
"github.com/degoke/health-ai-stack/pkg/store"
"github.com/degoke/health-ai-stack/pkg/types"
)
Expand Down Expand Up @@ -58,20 +59,22 @@ func (s *ResourceService) ProcessTransactionBundle(ctx context.Context, bundle *
}
committed = true

if err := s.syncDefinitionCatalog(ctx, writtenDefinitions, deletedDefinitions); err != nil {
return nil, exceptionErr("sync definition catalog from transaction bundle", err)
}

responseJSON, err := buildTransactionResponseBundle(responseEntries)
if err != nil {
return nil, exceptionErr("build transaction response bundle", err)
}

return &types.ResourceEnvelope{
response := &types.ResourceEnvelope{
ResourceType: "Bundle",
JSON: responseJSON,
Hash: mustHash(responseJSON),
}, nil
}
s.runPostCommit(ctx, hooks.ActionTransaction, response, nil)

if err := s.syncDefinitionCatalog(ctx, writtenDefinitions, deletedDefinitions); err != nil {
return nil, exceptionErr("sync definition catalog from transaction bundle", err)
}

return response, nil
}

type transactionBundle struct {
Expand Down Expand Up @@ -261,7 +264,7 @@ func (s *ResourceService) executeBundleCreate(
return bundleExecutionResult{}, conflictErr(fmt.Sprintf("resource already exists: %s/%s", envelope.ResourceType, id), nil)
}

written, err := s.applyWrite(ctx, session, envelope, store.VersionActionCreate)
written, err := s.applyWrite(ctx, session, envelope, store.VersionActionCreate, nil)
if err != nil {
return bundleExecutionResult{}, err
}
Expand Down Expand Up @@ -307,33 +310,29 @@ func (s *ResourceService) executeBundleUpdate(
}
}

exists, err := session.ResourceStore().Exists(ctx, resourceType, id)
previous, err := session.ResourceStore().Read(ctx, resourceType, id)
if err != nil {
return bundleExecutionResult{}, exceptionErr("check resource existence", err)
}
if !exists {
return bundleExecutionResult{}, notFoundErr(fmt.Sprintf("resource not found: %s/%s", resourceType, id), nil)
if isStoreNotFound(err) {
return bundleExecutionResult{}, notFoundErr(fmt.Sprintf("resource not found: %s/%s", resourceType, id), err)
}
return bundleExecutionResult{}, exceptionErr("read previous resource", err)
}
if entry.IfMatch != "" {
expected, ok := versionFromETag(entry.IfMatch)
if !ok {
return bundleExecutionResult{}, invalidErr("bundle ifMatch must contain one entity tag", nil)
}
current, err := session.ResourceStore().Read(ctx, resourceType, id)
if err != nil {
return bundleExecutionResult{}, exceptionErr("read current resource for ifMatch", err)
}
if expected != "*" && expected != current.VersionID {
if expected != "*" && expected != previous.VersionID {
return bundleExecutionResult{}, preconditionErr(fmt.Sprintf("resource version does not match expected version %q", expected), nil)
}
written, err := s.applyWriteExpectedVersion(ctx, session, envelope, store.VersionActionUpdate, expected)
written, err := s.applyWriteExpectedVersion(ctx, session, envelope, store.VersionActionUpdate, expected, hooks.ActionUpdate, previous)
if err != nil {
return bundleExecutionResult{}, err
}
return bundleExecutionFromWrite("200 OK", written), nil
}

written, err := s.applyWrite(ctx, session, envelope, store.VersionActionUpdate)
written, err := s.applyWrite(ctx, session, envelope, store.VersionActionUpdate, previous)
if err != nil {
return bundleExecutionResult{}, err
}
Expand Down
55 changes: 5 additions & 50 deletions pkg/core/conditional.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"context"
"fmt"

"github.com/degoke/health-ai-stack/pkg/hooks"
"github.com/degoke/health-ai-stack/pkg/store"
"github.com/degoke/health-ai-stack/pkg/types"
)
Expand Down Expand Up @@ -50,7 +51,7 @@ func (s *ResourceService) UpdateIfMatch(ctx context.Context, resource *types.Res
}
return nil, exceptionErr("read previous resource", err)
}
written, err := s.applyWriteExpectedVersion(ctx, session, envelope, store.VersionActionUpdate, expectedVersion)
written, err := s.applyWriteExpectedVersion(ctx, session, envelope, store.VersionActionUpdate, expectedVersion, hooks.ActionUpdate, previous)
if err != nil {
return nil, err
}
Expand All @@ -61,6 +62,7 @@ func (s *ResourceService) UpdateIfMatch(ctx context.Context, resource *types.Res
return nil, exceptionErr("commit write session", err)
}
committed = true
s.runPostCommit(ctx, hooks.ActionUpdate, written, previous)
return written, nil
}

Expand Down Expand Up @@ -97,6 +99,7 @@ func (s *ResourceService) DeleteIfMatch(ctx context.Context, resourceType, id, e
return exceptionErr("commit write session", err)
}
committed = true
s.runPostCommit(ctx, hooks.ActionDelete, current, current)
return nil
}

Expand All @@ -106,53 +109,5 @@ func (s *ResourceService) PatchIfMatch(ctx context.Context, resourceType, id str
if resourceType == "" || id == "" || expectedVersion == "" {
return nil, invalidErr("resourceType, id, and expected version are required", nil)
}
if len(patchJSON) == 0 {
return nil, invalidErr("patch body is required", nil)
}
session, err := s.sessions.BeginWrite(ctx)
if err != nil {
return nil, exceptionErr("begin write session", err)
}
committed := false
defer func() {
if !committed {
_ = session.Rollback(ctx)
}
}()
current, err := session.ResourceStore().Read(ctx, resourceType, id)
if err != nil {
if isStoreNotFound(err) {
return nil, notFoundErr(fmt.Sprintf("resource not found: %s/%s", resourceType, id), err)
}
return nil, exceptionErr("read resource for patch", err)
}
patchedJSON, err := applyPatchDocument(current.JSON, patchJSON)
if err != nil {
return nil, invalidErr("apply patch", err)
}
if err := validatePatchedIdentity(patchedJSON, resourceType, id); err != nil {
return nil, err
}
envelope := &types.ResourceEnvelope{ResourceType: resourceType, ID: id, JSON: patchedJSON}
envelope, err = s.normalizeEnvelope(envelope)
if err != nil {
return nil, err
}
if s.validator != nil {
if err := s.validator.ValidateResource(ctx, envelope); err != nil {
return nil, invalidErr("resource validation failed", err)
}
}
written, err := s.applyWriteExpectedVersion(ctx, session, envelope, store.VersionActionUpdate, expectedVersion)
if err != nil {
return nil, err
}
if err := s.removePreviousTerminology(ctx, session, current, written); err != nil {
return nil, exceptionErr("replace previous terminology projection", err)
}
if err := session.Commit(ctx); err != nil {
return nil, exceptionErr("commit write session", err)
}
committed = true
return written, nil
return s.patchAndCommit(ctx, resourceType, id, patchJSON, expectedVersion)
}
2 changes: 2 additions & 0 deletions pkg/core/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,8 @@
// - Indexer — search.Indexer invoked after resource/history persistence in session.
// - Outbox — sync.Outbox; when non-nil, events are appended via
// sync.WithWriteSession during each write (transactional through session EventStore).
// - Hooks — optional four-point intercept SPI; core runs pre-storage before persist
// and post-commit after a successful session commit.
//
// Methods:
//
Expand Down
Loading
Loading