324 lines
11 KiB
Go
324 lines
11 KiB
Go
package pipeline
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"chis_osi/contract"
|
|
"chis_osi/mapping"
|
|
"chis_osi/osi"
|
|
"chis_osi/source"
|
|
)
|
|
|
|
type Status string
|
|
|
|
const (
|
|
StatusDone Status = "done"
|
|
StatusRetry Status = "retry"
|
|
StatusFailed Status = "failed"
|
|
StatusManualReview Status = "manual_review"
|
|
)
|
|
|
|
type Action string
|
|
|
|
const (
|
|
ActionNone Action = "none"
|
|
ActionCreate Action = "create"
|
|
ActionUpdate Action = "update"
|
|
ActionSkip Action = "skip"
|
|
)
|
|
|
|
type Outcome struct {
|
|
Status Status `json:"status"`
|
|
Action Action `json:"action"`
|
|
Reason string `json:"reason,omitempty"`
|
|
ServiceID string `json:"serviceId,omitempty"`
|
|
ResponseCode string `json:"responseCode,omitempty"`
|
|
PHRIDHint string `json:"phrIdHint,omitempty"`
|
|
}
|
|
|
|
type UpsertEvent struct {
|
|
TraceID string
|
|
Status Status
|
|
Action Action
|
|
Reason string
|
|
ServiceID string
|
|
ResponseCode string
|
|
PHRIDHint string
|
|
}
|
|
|
|
type PHISStatusUpdate struct {
|
|
SourceRecordID string
|
|
BusinessID string
|
|
Status Status
|
|
Reason string
|
|
}
|
|
|
|
type UpdateTarget struct {
|
|
PHRID string
|
|
EMPIID string
|
|
CheckID string
|
|
IDCard string
|
|
ManaUnitID string
|
|
ManaDoctorID string
|
|
Status string
|
|
}
|
|
|
|
type Converter interface {
|
|
ConvertHealthRecord(source.HealthRecordTask) (contract.HealthRecordCreate, []mapping.ValidationError)
|
|
}
|
|
|
|
type ConverterFunc func(source.HealthRecordTask) (contract.HealthRecordCreate, []mapping.ValidationError)
|
|
|
|
func (f ConverterFunc) ConvertHealthRecord(task source.HealthRecordTask) (contract.HealthRecordCreate, []mapping.ValidationError) {
|
|
return f(task)
|
|
}
|
|
|
|
type HealthRecordClient interface {
|
|
FindHealthRecord(context.Context, osi.FindHealthRecordQuery) ([]contract.FindHealthRecord, osi.Result, error)
|
|
CreateHealthRecord(context.Context, contract.HealthRecordCreate) (contract.HealthRecordSaveResult, osi.Result, error)
|
|
UpdateHealthRecord(context.Context, contract.HealthRecordCreate) (contract.HealthRecordSaveResult, osi.Result, error)
|
|
}
|
|
|
|
type IdempotencyStore interface {
|
|
Acquire(context.Context, string) (IdempotencyLease, error)
|
|
Complete(context.Context, string, string) error
|
|
Release(context.Context, string, string) error
|
|
}
|
|
|
|
type IdempotencyState string
|
|
|
|
const (
|
|
IdempotencyAcquired IdempotencyState = "acquired"
|
|
IdempotencyCompleted IdempotencyState = "completed"
|
|
IdempotencyInProgress IdempotencyState = "in_progress"
|
|
)
|
|
|
|
type IdempotencyLease struct {
|
|
State IdempotencyState
|
|
Token string
|
|
}
|
|
|
|
type NotificationError struct {
|
|
Event *UpsertEvent
|
|
Update *PHISStatusUpdate
|
|
PriorErr error
|
|
ReportErr error
|
|
StatusErr error
|
|
}
|
|
|
|
func (e *NotificationError) Error() string {
|
|
return "finish upsert side effects: " + e.Unwrap().Error()
|
|
}
|
|
func (e *NotificationError) Unwrap() error { return errors.Join(e.PriorErr, e.ReportErr, e.StatusErr) }
|
|
|
|
type ReportSink interface {
|
|
Publish(context.Context, UpsertEvent) error
|
|
}
|
|
|
|
type PHISStatusWriter interface {
|
|
WriteStatus(context.Context, PHISStatusUpdate) error
|
|
}
|
|
|
|
type Dependencies struct {
|
|
Converter Converter
|
|
Client HealthRecordClient
|
|
Idempotency IdempotencyStore
|
|
Reports ReportSink
|
|
PHISStatuses PHISStatusWriter
|
|
}
|
|
|
|
type HealthRecordUpsertService struct {
|
|
deps Dependencies
|
|
}
|
|
|
|
const idempotencyCleanupTimeout = 5 * time.Second
|
|
|
|
func NewHealthRecordUpsertService(deps Dependencies) *HealthRecordUpsertService {
|
|
return &HealthRecordUpsertService{deps: deps}
|
|
}
|
|
|
|
func (s *HealthRecordUpsertService) Upsert(ctx context.Context, task source.HealthRecordTask) (Outcome, error) {
|
|
if err := s.validateDependencies(); err != nil {
|
|
return Outcome{}, err
|
|
}
|
|
|
|
req, validationErrors := s.deps.Converter.ConvertHealthRecord(task)
|
|
if len(validationErrors) > 0 {
|
|
return s.finish(ctx, task, Outcome{Status: StatusFailed, Action: ActionNone, Reason: "mapping validation failed: " + validationErrors[0].Field})
|
|
}
|
|
|
|
idempotencyKey := healthRecordIdempotencyKey(req.HealthRecord.CheckID, req.HealthRecord.IDCard)
|
|
lease, err := s.deps.Idempotency.Acquire(ctx, idempotencyKey)
|
|
if err != nil {
|
|
return s.finish(ctx, task, Outcome{Status: StatusRetry, Action: ActionNone, Reason: "idempotency lookup failed"}, err)
|
|
}
|
|
if lease.State == IdempotencyCompleted {
|
|
return s.finish(ctx, task, Outcome{Status: StatusDone, Action: ActionSkip, Reason: "already completed"})
|
|
}
|
|
if lease.State == IdempotencyInProgress {
|
|
return s.finish(ctx, task, Outcome{Status: StatusRetry, Action: ActionSkip, Reason: "another delivery is in progress"})
|
|
}
|
|
if lease.State != IdempotencyAcquired || strings.TrimSpace(lease.Token) == "" {
|
|
return s.finish(ctx, task, Outcome{Status: StatusRetry, Action: ActionNone, Reason: "invalid idempotency lease"})
|
|
}
|
|
|
|
records, queryResult, err := s.deps.Client.FindHealthRecord(ctx, osi.FindHealthRecordQuery{IDCard: req.HealthRecord.IDCard})
|
|
if err != nil || !queryResult.Success {
|
|
status := classifyFailure(queryResult)
|
|
return s.finishAcquired(ctx, task, idempotencyKey, lease.Token, Outcome{Status: status, Action: ActionNone, Reason: "CHIS query failed", ServiceID: osi.ServiceIDJKDAFind, ResponseCode: queryResult.Code})
|
|
}
|
|
if records == nil {
|
|
return s.finishAcquired(ctx, task, idempotencyKey, lease.Token, Outcome{Status: StatusFailed, Action: ActionNone, Reason: "CHIS query returned null or missing data", ServiceID: osi.ServiceIDJKDAFind, ResponseCode: queryResult.Code})
|
|
}
|
|
|
|
if len(records) > 1 {
|
|
return s.finishAcquired(ctx, task, idempotencyKey, lease.Token, Outcome{Status: StatusManualReview, Action: ActionNone, Reason: "multiple CHIS health records found", ServiceID: osi.ServiceIDJKDAFind, ResponseCode: queryResult.Code})
|
|
}
|
|
if len(records) == 0 {
|
|
return s.create(ctx, task, idempotencyKey, lease.Token, req)
|
|
}
|
|
|
|
target := newUpdateTarget(records[0].HealthRecord)
|
|
if reason := validateUpdateTarget(target, req); reason != "" {
|
|
return s.finishAcquired(ctx, task, idempotencyKey, lease.Token, Outcome{Status: StatusManualReview, Action: ActionNone, Reason: reason, ServiceID: osi.ServiceIDJKDAFind, ResponseCode: queryResult.Code, PHRIDHint: maskIdentifier(target.PHRID)})
|
|
}
|
|
req.BaseInfo.PHRID = target.PHRID
|
|
req.HealthRecord.PhrID = target.PHRID
|
|
return s.update(ctx, task, idempotencyKey, lease.Token, req)
|
|
}
|
|
|
|
func (s *HealthRecordUpsertService) create(ctx context.Context, task source.HealthRecordTask, key, token string, req contract.HealthRecordCreate) (Outcome, error) {
|
|
saved, result, err := s.deps.Client.CreateHealthRecord(ctx, req)
|
|
outcome := Outcome{Status: StatusDone, Action: ActionCreate, ServiceID: osi.ServiceIDJKDACreate, ResponseCode: result.Code, PHRIDHint: maskIdentifier(saved.PhrID)}
|
|
if err != nil || !result.Success {
|
|
outcome.Status = classifyFailure(result)
|
|
outcome.Reason = "CHIS create failed"
|
|
return s.finishAcquired(ctx, task, key, token, outcome)
|
|
}
|
|
return s.complete(ctx, task, key, token, outcome)
|
|
}
|
|
|
|
func (s *HealthRecordUpsertService) update(ctx context.Context, task source.HealthRecordTask, key, token string, req contract.HealthRecordCreate) (Outcome, error) {
|
|
saved, result, err := s.deps.Client.UpdateHealthRecord(ctx, req)
|
|
phrID := saved.PhrID
|
|
if phrID == "" {
|
|
phrID = req.HealthRecord.PhrID
|
|
}
|
|
outcome := Outcome{Status: StatusDone, Action: ActionUpdate, ServiceID: osi.ServiceIDJKDAUpdate, ResponseCode: result.Code, PHRIDHint: maskIdentifier(phrID)}
|
|
if err != nil || !result.Success {
|
|
outcome.Status = classifyFailure(result)
|
|
outcome.Reason = "CHIS update failed"
|
|
return s.finishAcquired(ctx, task, key, token, outcome)
|
|
}
|
|
return s.complete(ctx, task, key, token, outcome)
|
|
}
|
|
|
|
func (s *HealthRecordUpsertService) complete(ctx context.Context, task source.HealthRecordTask, key, token string, outcome Outcome) (Outcome, error) {
|
|
cleanupCtx, cancel := idempotencyCleanupContext(ctx)
|
|
defer cancel()
|
|
if err := s.deps.Idempotency.Complete(cleanupCtx, key, token); err != nil {
|
|
return outcome, fmt.Errorf("persist successful CHIS write before notification: %w", err)
|
|
}
|
|
return s.finish(ctx, task, outcome)
|
|
}
|
|
|
|
func (s *HealthRecordUpsertService) finishAcquired(ctx context.Context, task source.HealthRecordTask, key, token string, outcome Outcome) (Outcome, error) {
|
|
cleanupCtx, cancel := idempotencyCleanupContext(ctx)
|
|
defer cancel()
|
|
return s.finish(ctx, task, outcome, s.deps.Idempotency.Release(cleanupCtx, key, token))
|
|
}
|
|
|
|
func idempotencyCleanupContext(ctx context.Context) (context.Context, context.CancelFunc) {
|
|
return context.WithTimeout(context.WithoutCancel(ctx), idempotencyCleanupTimeout)
|
|
}
|
|
|
|
func (s *HealthRecordUpsertService) finish(ctx context.Context, task source.HealthRecordTask, outcome Outcome, prior ...error) (Outcome, error) {
|
|
event := UpsertEvent{
|
|
TraceID: traceID(task.SourceRecordID(), task.BusinessID), Status: outcome.Status,
|
|
Action: outcome.Action, Reason: outcome.Reason, ServiceID: outcome.ServiceID,
|
|
ResponseCode: outcome.ResponseCode, PHRIDHint: outcome.PHRIDHint,
|
|
}
|
|
status := PHISStatusUpdate{SourceRecordID: task.SourceRecordID(), BusinessID: task.BusinessID, Status: outcome.Status, Reason: outcome.Reason}
|
|
priorErr := errors.Join(prior...)
|
|
reportErr := s.deps.Reports.Publish(ctx, event)
|
|
statusErr := s.deps.PHISStatuses.WriteStatus(ctx, status)
|
|
if priorErr != nil || reportErr != nil || statusErr != nil {
|
|
notificationErr := &NotificationError{PriorErr: priorErr, ReportErr: reportErr, StatusErr: statusErr}
|
|
if reportErr != nil {
|
|
notificationErr.Event = &event
|
|
}
|
|
if statusErr != nil {
|
|
notificationErr.Update = &status
|
|
}
|
|
return outcome, notificationErr
|
|
}
|
|
return outcome, nil
|
|
}
|
|
|
|
func (s *HealthRecordUpsertService) validateDependencies() error {
|
|
if s == nil || s.deps.Converter == nil || s.deps.Client == nil || s.deps.Idempotency == nil || s.deps.Reports == nil || s.deps.PHISStatuses == nil {
|
|
return fmt.Errorf("health record upsert dependencies are incomplete")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func classifyFailure(result osi.Result) Status {
|
|
if result.Retryable {
|
|
return StatusRetry
|
|
}
|
|
return StatusFailed
|
|
}
|
|
|
|
func newUpdateTarget(record contract.HealthRecord) UpdateTarget {
|
|
checkID := ""
|
|
if record.CheckID != nil {
|
|
checkID = *record.CheckID
|
|
}
|
|
return UpdateTarget{
|
|
PHRID: record.PhrID, EMPIID: record.EmpiID, CheckID: checkID, IDCard: record.IDCard,
|
|
ManaUnitID: record.ManaUnitID, ManaDoctorID: record.ManaDoctorID, Status: record.Status,
|
|
}
|
|
}
|
|
|
|
func validateUpdateTarget(target UpdateTarget, req contract.HealthRecordCreate) string {
|
|
if strings.TrimSpace(target.IDCard) != strings.TrimSpace(req.HealthRecord.IDCard) {
|
|
return "CHIS identity does not exactly match source"
|
|
}
|
|
if strings.TrimSpace(target.PHRID) == "" {
|
|
return "CHIS record is missing phrId"
|
|
}
|
|
if strings.TrimSpace(target.Status) != "1" {
|
|
return "CHIS record status is not updateable"
|
|
}
|
|
if strings.TrimSpace(target.ManaUnitID) != strings.TrimSpace(req.HealthRecord.ManaUnitID) {
|
|
return "CHIS record belongs to another organization"
|
|
}
|
|
if strings.TrimSpace(target.ManaDoctorID) != strings.TrimSpace(req.HealthRecord.ManaDoctorID) {
|
|
return "CHIS record belongs to another doctor"
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func healthRecordIdempotencyKey(checkID, idCard string) string {
|
|
sum := sha256.Sum256([]byte(strings.TrimSpace(checkID) + "|" + strings.TrimSpace(idCard)))
|
|
return fmt.Sprintf("health-record:%x", sum[:])
|
|
}
|
|
|
|
func traceID(sourceRecordID, businessID string) string {
|
|
sum := sha256.Sum256([]byte(strings.TrimSpace(sourceRecordID) + "|" + strings.TrimSpace(businessID)))
|
|
return fmt.Sprintf("hr-%x", sum[:8])
|
|
}
|
|
|
|
func maskIdentifier(value string) string {
|
|
value = strings.TrimSpace(value)
|
|
if len(value) <= 4 {
|
|
return "****"
|
|
}
|
|
return "****" + value[len(value)-4:]
|
|
}
|