ReadRelationships reads relation tuples from the storage based on the given filter and pagination.
(ctx context.Context, tenantID string, filter *base.TupleFilter, snap string, pagination database.Pagination)
| 116 | |
| 117 | // ReadRelationships reads relation tuples from the storage based on the given filter and pagination. |
| 118 | func (r *DataReader) ReadRelationships(ctx context.Context, tenantID string, filter *base.TupleFilter, snap string, pagination database.Pagination) (collection *database.TupleCollection, ct database.EncodedContinuousToken, err error) { |
| 119 | // Start a new trace span and end it when the function exits. |
| 120 | ctx, span := internal.Tracer.Start(ctx, "data-reader.read-relationships") |
| 121 | defer span.End() |
| 122 | // Log read operation |
| 123 | slog.DebugContext(ctx, "reading relationships for tenant_id", slog.String("tenant_id", tenantID)) |
| 124 | // Decode snapshot token |
| 125 | // Decode the snapshot value. |
| 126 | var st token.SnapToken |
| 127 | st, err = snapshot.EncodedToken{Value: snap}.Decode() |
| 128 | if err != nil { |
| 129 | return nil, nil, utils.HandleError(ctx, span, err, base.ErrorCode_ERROR_CODE_INTERNAL) |
| 130 | } |
| 131 | |
| 132 | // Build the relationships query based on the provided filter, snapshot value, and pagination settings. |
| 133 | builder := r.database.Builder.Select("id, entity_type, entity_id, relation, subject_type, subject_id, subject_relation").From(RelationTuplesTable).Where(squirrel.Eq{"tenant_id": tenantID}) |
| 134 | builder = utils.TuplesFilterQueryForSelectBuilder(builder, filter) |
| 135 | builder = utils.SnapshotQuery(builder, st.(snapshot.Token).Value.Uint, st.(snapshot.Token).Snapshot) |
| 136 | |
| 137 | // Apply the pagination token and limit to the query. |
| 138 | if pagination.Token() != "" { |
| 139 | var t database.ContinuousToken |
| 140 | t, err = utils.EncodedContinuousToken{Value: pagination.Token()}.Decode() |
| 141 | if err != nil { |
| 142 | return nil, nil, utils.HandleError(ctx, span, err, base.ErrorCode_ERROR_CODE_INVALID_CONTINUOUS_TOKEN) |
| 143 | } |
| 144 | var v uint64 |
| 145 | v, err = strconv.ParseUint(t.(utils.ContinuousToken).Value, 10, 64) |
| 146 | if err != nil { |
| 147 | return nil, nil, utils.HandleError(ctx, span, err, base.ErrorCode_ERROR_CODE_INVALID_CONTINUOUS_TOKEN) |
| 148 | } |
| 149 | builder = builder.Where(squirrel.GtOrEq{"id": v}) |
| 150 | } |
| 151 | |
| 152 | builder = builder.OrderBy("id") |
| 153 | |
| 154 | if pagination.PageSize() != 0 { |
| 155 | builder = builder.Limit(uint64(pagination.PageSize() + 1)) |
| 156 | } |
| 157 | |
| 158 | // Generate the SQL query and arguments. |
| 159 | var query string |
| 160 | var args []interface{} |
| 161 | query, args, err = builder.ToSql() |
| 162 | if err != nil { |
| 163 | return nil, database.NewNoopContinuousToken().Encode(), utils.HandleError(ctx, span, err, base.ErrorCode_ERROR_CODE_SQL_BUILDER) |
| 164 | } |
| 165 | // Log generated query |
| 166 | slog.DebugContext(ctx, "generated sql query", slog.String("query", query), "with args", slog.Any("arguments", args)) |
| 167 | // Execute query |
| 168 | // Execute the query and retrieve the rows. |
| 169 | var rows pgx.Rows |
| 170 | rows, err = r.database.ReadPool.Query(ctx, query, args...) |
| 171 | if err != nil { |
| 172 | return nil, database.NewNoopContinuousToken().Encode(), utils.HandleError(ctx, span, err, base.ErrorCode_ERROR_CODE_EXECUTION) |
| 173 | } |
| 174 | defer rows.Close() |
| 175 |
nothing calls this directly
no test coverage detected