package httpapi import ( "context" "encoding/json" "fmt" "strings" "time" "github.com/dtoro/oikos/internal/db/sqlcgen" "github.com/dtoro/oikos/internal/domain" "github.com/dtoro/oikos/internal/httpapi/gen" "github.com/dtoro/oikos/internal/observability" ) // ─── Relationships ───────────────────────────────────────────────────── func (s *Server) CreateRelationship(ctx context.Context, req gen.CreateRelationshipRequestObject) (gen.CreateRelationshipResponseObject, error) { if req.Body == nil { return nil, fmt.Errorf("%w: request body is required", domain.ErrInvalidInput) } sourceID, err := s.resolveEntityID(ctx, req.Body.Source) if err != nil { return nil, err } targetID, err := s.resolveEntityID(ctx, req.Body.Target) if err != nil { return nil, err } attrsJSON := []byte("{}") if req.Body.Attributes != nil { attrsJSON, _ = json.Marshal(req.Body.Attributes) } tx, err := s.pool.Begin(ctx) if err != nil { return nil, err } defer tx.Rollback(ctx) _, err = tx.Exec(ctx, ` INSERT INTO relationships (source_id, target_id, type, attributes, valid_from) VALUES ($1, $2, $3, $4, now())`, sourceID, targetID, req.Body.Type, attrsJSON) if err != nil { if strings.Contains(err.Error(), "unique") || strings.Contains(err.Error(), "duplicate") { return nil, fmt.Errorf("%w: relationship %s:%s:%s already exists", domain.ErrAlreadyExists, req.Body.Source, req.Body.Type, req.Body.Target) } return nil, err } rel := gen.Relationship{ Source: req.Body.Source, Target: req.Body.Target, Type: req.Body.Type, ValidFrom: time.Now(), } if req.Body.Attributes != nil { rel.Attributes = req.Body.Attributes } actorType, actor := actorInfo(ctx) if auditErr := observability.Audit(ctx, sqlcgen.New(tx), actorType, actor, "create", nil, "POST", "/api/v1/relationships", "", map[string]any{"source": req.Body.Source, "target": req.Body.Target, "type": req.Body.Type}); auditErr != nil { return nil, auditErr } if err := tx.Commit(ctx); err != nil { return nil, err } return gen.CreateRelationship201JSONResponse(rel), nil } func (s *Server) EndRelationship(ctx context.Context, req gen.EndRelationshipRequestObject) (gen.EndRelationshipResponseObject, error) { sourceID, err := s.resolveEntityID(ctx, req.Params.Source) if err != nil { return nil, err } targetID, err := s.resolveEntityID(ctx, req.Params.Target) if err != nil { return nil, err } tx, err := s.pool.Begin(ctx) if err != nil { return nil, err } defer tx.Rollback(ctx) result, err := sqlcgen.New(tx).EndCurrentRelationship(ctx, sqlcgen.EndCurrentRelationshipParams{ SourceID: sourceID, TargetID: targetID, Type: req.Params.RelType, }) if err != nil { return nil, err } if result == 0 { return nil, fmt.Errorf("%w: active relationship %s:%s:%s", domain.ErrNotFound, req.Params.Source, req.Params.RelType, req.Params.Target) } actorType, actor := actorInfo(ctx) if auditErr := observability.Audit(ctx, sqlcgen.New(tx), actorType, actor, "delete", nil, "DELETE", "/api/v1/relationships", "", map[string]any{"source": req.Params.Source, "target": req.Params.Target, "type": req.Params.RelType}); auditErr != nil { return nil, auditErr } if err := tx.Commit(ctx); err != nil { return nil, err } return gen.EndRelationship204Response{}, nil }