( input: BuildAgentGraphTraceSnapshotInput, )
| 99 | } |
| 100 | |
| 101 | export function buildAgentGraphTraceSnapshot( |
| 102 | input: BuildAgentGraphTraceSnapshotInput, |
| 103 | ): AgentGraphTraceSnapshot { |
| 104 | const validated = validateTopology(input.topology); |
| 105 | const topologyFingerprint = fingerprintTopology(input.topology.graphId, validated); |
| 106 | const replay = |
| 107 | input.records.length > 0 |
| 108 | ? replayAgentGraphRecords(input.records) |
| 109 | : { |
| 110 | graphId: input.topology.graphId, |
| 111 | appliedRecordIds: [], |
| 112 | operators: {}, |
| 113 | }; |
| 114 | |
| 115 | if (replay.graphId !== input.topology.graphId) { |
| 116 | throw new Error( |
| 117 | `Trace topology ${input.topology.graphId} cannot observe records from graph ${replay.graphId}`, |
| 118 | ); |
| 119 | } |
| 120 | |
| 121 | const recordsById = new Map(input.records.map((record) => [record.recordId, record])); |
| 122 | const orderedRecords = replay.appliedRecordIds.map((recordId) => recordsById.get(recordId)!); |
| 123 | for (const record of orderedRecords) { |
| 124 | const binding = validated.operatorsById.get(record.operatorId); |
| 125 | if (!binding) { |
| 126 | throw new Error( |
| 127 | `Graph record ${record.recordId} references unknown topology operator ${record.operatorId}`, |
| 128 | ); |
| 129 | } |
| 130 | if (binding.sessionId !== record.sessionId) { |
| 131 | throw new Error( |
| 132 | `Topology operator ${record.operatorId} is bound to ${binding.sessionId}, record uses ${record.sessionId}`, |
| 133 | ); |
| 134 | } |
| 135 | } |
| 136 | |
| 137 | const topologicalIndex = new Map( |
| 138 | validated.topologicalOrder.map((operatorId, index) => [operatorId, index]), |
| 139 | ); |
| 140 | const compareOperators = (a: string, b: string): number => |
| 141 | topologicalIndex.get(a)! - topologicalIndex.get(b)! || compareAgentGraphIdentity(a, b); |
| 142 | |
| 143 | const replayOperators = new Map(Object.entries(replay.operators)); |
| 144 | const operators = new Map<string, AgentGraphTraceOperatorState>(); |
| 145 | for (const operatorId of validated.topologicalOrder) { |
| 146 | const binding = validated.operatorsById.get(operatorId)!; |
| 147 | const runtimeState = replayOperators.get(operatorId); |
| 148 | operators.set(operatorId, { |
| 149 | operatorId, |
| 150 | sessionId: binding.sessionId, |
| 151 | topologicalIndex: topologicalIndex.get(operatorId)!, |
| 152 | upstreamOperatorIds: uniqueOperatorIds( |
| 153 | (validated.incoming.get(operatorId) ?? []).map((edge) => edge.fromOperatorId), |
| 154 | compareOperators, |
| 155 | ), |
| 156 | downstreamOperatorIds: uniqueOperatorIds( |
| 157 | (validated.outgoing.get(operatorId) ?? []).map((edge) => edge.toOperatorId), |
| 158 | compareOperators, |
no test coverage detected