MCPcopy Create free account
hub / github.com/apache/maka / buildAgentGraphTraceSnapshot

Function buildAgentGraphTraceSnapshot

packages/runtime/src/stream-graph-trace.ts:101–218  ·  view source on GitHub ↗
(
  input: BuildAgentGraphTraceSnapshotInput,
)

Source from the content-addressed store, hash-verified

99}
100
101export 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,

Calls 12

validateTopologyFunction · 0.85
fingerprintTopologyFunction · 0.85
replayAgentGraphRecordsFunction · 0.85
uniqueOperatorIdsFunction · 0.85
cloneOperatorStateFunction · 0.85
compareOperatorsFunction · 0.85
traceRouteIdFunction · 0.85
getMethod · 0.65
entriesMethod · 0.65
setMethod · 0.65
pushMethod · 0.65

Tested by

no test coverage detected