Skip to content

Commit 5dc6523

Browse files
[SDFAB-189] Add UPF programmable behaviour (#276)
* Initial move of upf programmable * Cleanup tests for upf programmable * Fix checkstyle * Fix pom file * Add UPF behaviour when pipeconf requires it Also, use do not manage UE limit on physical pipeline. * Force to use local artifacs * Fix read all counters usage and new tests * Update full profile suffix. * Make mocks more generic and independent from upf constants * Add dependency to fabric v1model and remove redundant interfaces * Add missing license headers * Remove dependency to fabric-v1model * Use different name for distributed structures * Rename upf store interface name * Build using snapshots * Fix MockReadRequest * Fix failure to find onos-build-conf Co-authored-by: Carmelo Cascone <carmelo@opennetworking.org>
1 parent bcbb0e7 commit 5dc6523

27 files changed

Lines changed: 3538 additions & 9 deletions

pom.xml

Lines changed: 19 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ SPDX-License-Identifier: LicenseRef-ONF-Member-Only-1.0
1111
<parent>
1212
<groupId>org.onosproject</groupId>
1313
<artifactId>onos-dependencies</artifactId>
14-
<version>2.5.2-b1</version>
14+
<version>2.5.2-SNAPSHOT</version>
1515
</parent>
1616

1717
<groupId>org.stratumproject</groupId>
@@ -40,22 +40,22 @@ SPDX-License-Identifier: LicenseRef-ONF-Member-Only-1.0
4040
</dependency>
4141

4242
<dependency>
43-
<groupId>${trellis.api.groupId}</groupId>
44-
<artifactId>${trellis.api.artifactId}</artifactId>
45-
<version>${trellis.api.version}</version>
43+
<groupId>org.onosproject</groupId>
44+
<artifactId>onos-drivers-p4runtime</artifactId>
45+
<version>${onos.version}</version>
4646
<scope>provided</scope>
4747
</dependency>
4848

4949
<dependency>
50-
<groupId>org.onosproject</groupId>
51-
<artifactId>onlab-osgi</artifactId>
52-
<version>${onos.version}</version>
50+
<groupId>${trellis.api.groupId}</groupId>
51+
<artifactId>${trellis.api.artifactId}</artifactId>
52+
<version>${trellis.api.version}</version>
5353
<scope>provided</scope>
5454
</dependency>
5555

5656
<dependency>
5757
<groupId>org.onosproject</groupId>
58-
<artifactId>onos-pipelines-fabric-api</artifactId>
58+
<artifactId>onlab-osgi</artifactId>
5959
<version>${onos.version}</version>
6060
<scope>provided</scope>
6161
</dependency>
@@ -181,6 +181,17 @@ SPDX-License-Identifier: LicenseRef-ONF-Member-Only-1.0
181181
</snapshots>
182182
</repository>
183183
</repositories>
184+
<pluginRepositories>
185+
<pluginRepository>
186+
<id>snapshots</id>
187+
<url>https://oss.sonatype.org/content/repositories/snapshots</url>
188+
<snapshots>
189+
<enabled>true</enabled>
190+
<updatePolicy>always</updatePolicy>
191+
<checksumPolicy>fail</checksumPolicy>
192+
</snapshots>
193+
</pluginRepository>
194+
</pluginRepositories>
184195

185196
<distributionManagement>
186197
<snapshotRepository>

src/main/java/org/stratumproject/fabric/tna/PipeconfLoader.java

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55

66
import org.onosproject.core.CoreService;
77
import org.onosproject.net.behaviour.Pipeliner;
8+
import org.onosproject.net.behaviour.upf.UpfProgrammable;
89
import org.onosproject.net.behaviour.inbandtelemetry.IntProgrammable;
910
import org.onosproject.net.pi.model.DefaultPiPipeconf;
1011
import org.onosproject.net.pi.model.PiPipeconf;
@@ -26,6 +27,7 @@
2627
import org.stratumproject.fabric.tna.behaviour.FabricIntProgrammable;
2728
import org.stratumproject.fabric.tna.behaviour.FabricInterpreter;
2829
import org.stratumproject.fabric.tna.behaviour.pipeliner.FabricPipeliner;
30+
import org.stratumproject.fabric.tna.behaviour.upf.FabricUpfProgrammable;
2931

3032
import java.io.File;
3133
import java.io.FileNotFoundException;
@@ -68,7 +70,8 @@ public class PipeconfLoader {
6870
private static final String PIPELINE_CONFIG = "pipeline_config.pb.bin";
6971

7072
private static final String INT_PROFILE_SUFFIX = "-int";
71-
private static final String FULL_PROFILE_SUFFIX = "-full";
73+
private static final String UPF_PROFILE_SUFFIX = "-spgw";
74+
private static final String FULL_PROFILE_SUFFIX = "-spgw-int";
7275

7376
@Activate
7477
public void activate() {
@@ -132,6 +135,13 @@ private PiPipeconf buildPipeconfFromPath(String path) {
132135
builder.addBehaviour(IntProgrammable.class, FabricIntProgrammable.class);
133136
}
134137

138+
// Add UpfProgrammable behaviour for UPF-enabled profiles.
139+
if (profile.endsWith(UPF_PROFILE_SUFFIX) ||
140+
profile.endsWith(FULL_PROFILE_SUFFIX)) {
141+
builder.addBehaviour(UpfProgrammable.class, FabricUpfProgrammable.class);
142+
}
143+
144+
135145
return builder.build();
136146
}
137147

src/main/java/org/stratumproject/fabric/tna/behaviour/Constants.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,11 @@ public final class Constants {
3535
public static final int DEFAULT_PW_TRANSPORT_VLAN = 4090;
3636
public static final int PKT_IN_MIRROR_SESSION_ID = 0x210;
3737

38+
// UPF related constants
39+
public static final int UPF_INTERFACE_ACCESS = 1;
40+
public static final int UPF_INTERFACE_CORE = 2;
41+
public static final int UPF_INTERFACE_DBUF = 3;
42+
3843
// hide default constructor
3944
private Constants() {
4045
}

src/main/java/org/stratumproject/fabric/tna/behaviour/FabricCapabilities.java

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
import static com.google.common.base.Preconditions.checkNotNull;
1616
import static org.onosproject.net.pi.model.PiPipeconf.ExtensionType.CPU_PORT_TXT;
1717
import static org.slf4j.LoggerFactory.getLogger;
18+
import static org.stratumproject.fabric.tna.behaviour.P4InfoConstants.FABRIC_INGRESS_SPGW_DOWNLINK_PDRS;
1819

1920
/**
2021
* Representation of the capabilities of a given fabric-tna pipeconf.
@@ -83,6 +84,17 @@ public Optional<Integer> cpuPort() {
8384
}
8485
}
8586

87+
/**
88+
* Returns true if the pipeconf supports UPF capabilities, false otherwise.
89+
*
90+
* @return boolean
91+
*/
92+
public boolean supportUpf() {
93+
return pipeconf.pipelineModel()
94+
.table(FABRIC_INGRESS_SPGW_DOWNLINK_PDRS)
95+
.isPresent();
96+
}
97+
8698
public boolean supportDoubleVlanTerm() {
8799
// TODO: re-enable support for double-vlan
88100
// FIXME: next_vlan has been moved to pre_next
Lines changed: 234 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,234 @@
1+
// Copyright 2020-present Open Networking Foundation
2+
// SPDX-License-Identifier: LicenseRef-ONF-Member-Only-1.0
3+
package org.stratumproject.fabric.tna.behaviour.upf;
4+
5+
import com.google.common.collect.BiMap;
6+
import com.google.common.collect.ImmutableBiMap;
7+
import com.google.common.collect.Maps;
8+
import org.onosproject.net.behaviour.upf.PacketDetectionRule;
9+
import org.onlab.packet.Ip4Address;
10+
import org.onlab.util.ImmutableByteSequence;
11+
import org.onlab.util.KryoNamespace;
12+
import org.onosproject.store.serializers.KryoNamespaces;
13+
import org.onosproject.store.service.ConsistentMap;
14+
import org.onosproject.store.service.DistributedSet;
15+
import org.onosproject.store.service.MapEvent;
16+
import org.onosproject.store.service.MapEventListener;
17+
import org.onosproject.store.service.Serializer;
18+
import org.onosproject.store.service.StorageService;
19+
import org.osgi.service.component.annotations.Activate;
20+
import org.osgi.service.component.annotations.Component;
21+
import org.osgi.service.component.annotations.Deactivate;
22+
import org.osgi.service.component.annotations.Reference;
23+
import org.osgi.service.component.annotations.ReferenceCardinality;
24+
import org.slf4j.Logger;
25+
import org.slf4j.LoggerFactory;
26+
27+
import java.util.HashSet;
28+
import java.util.Map;
29+
import java.util.Objects;
30+
import java.util.Set;
31+
32+
import static com.google.common.base.Preconditions.checkNotNull;
33+
34+
/**
35+
* Distributed implementation of FabricUpfStore.
36+
*/
37+
// FIXME: this store is generic and not tied to a single device, should we have a store based on deviceId?
38+
@Component(immediate = true, service = DistributedFabricUpfStore.class)
39+
public final class DistributedFabricUpfStore implements FabricUpfStore {
40+
41+
private final Logger log = LoggerFactory.getLogger(getClass());
42+
43+
@Reference(cardinality = ReferenceCardinality.MANDATORY)
44+
protected StorageService storageService;
45+
46+
protected static final String FAR_ID_MAP_NAME = "fabric-upf-far-id-tna";
47+
protected static final String BUFFER_FAR_ID_SET_NAME = "fabric-upf-buffer-far-id-tna";
48+
protected static final String FAR_ID_UE_MAP_NAME = "fabric-upf-far-id-ue-tna";
49+
protected static final KryoNamespace.Builder SERIALIZER = KryoNamespace.newBuilder()
50+
.register(KryoNamespaces.API)
51+
.register(UpfRuleIdentifier.class);
52+
53+
// Mapping between scheduling priority ranges with Tofino priority queues
54+
// i.e., default queues are 8 in Tofino
55+
private static final BiMap<Integer, Integer> SCHEDULING_PRIORITY_MAP
56+
= new ImmutableBiMap.Builder<Integer, Integer>()
57+
// Highest scheduling priority for 3GPP is 1 and highest Tofino queue priority is 7
58+
.put(1, 5)
59+
.put(6, 4)
60+
.put(7, 3)
61+
.put(8, 2)
62+
.put(9, 1)
63+
.build();
64+
65+
// Distributed local FAR ID to global FAR ID mapping
66+
protected ConsistentMap<UpfRuleIdentifier, Integer> farIdMap;
67+
private MapEventListener<UpfRuleIdentifier, Integer> farIdMapListener;
68+
// Local, reversed copy of farIdMapper for better reverse lookup performance
69+
protected Map<Integer, UpfRuleIdentifier> reverseFarIdMap;
70+
private int nextGlobalFarId = 1;
71+
72+
protected DistributedSet<UpfRuleIdentifier> bufferFarIds;
73+
protected ConsistentMap<UpfRuleIdentifier, Set<Ip4Address>> farIdToUeAddrs;
74+
75+
@Activate
76+
protected void activate() {
77+
// Allow unit test to inject farIdMap here.
78+
if (storageService != null) {
79+
this.farIdMap = storageService.<UpfRuleIdentifier, Integer>consistentMapBuilder()
80+
.withName(FAR_ID_MAP_NAME)
81+
.withRelaxedReadConsistency()
82+
.withSerializer(Serializer.using(SERIALIZER.build()))
83+
.build();
84+
this.bufferFarIds = storageService.<UpfRuleIdentifier>setBuilder()
85+
.withName(BUFFER_FAR_ID_SET_NAME)
86+
.withRelaxedReadConsistency()
87+
.withSerializer(Serializer.using(SERIALIZER.build()))
88+
.build().asDistributedSet();
89+
this.farIdToUeAddrs = storageService.<UpfRuleIdentifier, Set<Ip4Address>>consistentMapBuilder()
90+
.withName(FAR_ID_UE_MAP_NAME)
91+
.withRelaxedReadConsistency()
92+
.withSerializer(Serializer.using(SERIALIZER.build()))
93+
.build();
94+
95+
}
96+
farIdMapListener = new FarIdMapListener();
97+
farIdMap.addListener(farIdMapListener);
98+
99+
reverseFarIdMap = Maps.newHashMap();
100+
farIdMap.entrySet().forEach(entry -> reverseFarIdMap.put(entry.getValue().value(), entry.getKey()));
101+
102+
log.info("Started");
103+
}
104+
105+
@Deactivate
106+
protected void deactivate() {
107+
farIdMap.removeListener(farIdMapListener);
108+
farIdMap.destroy();
109+
reverseFarIdMap.clear();
110+
111+
log.info("Stopped");
112+
}
113+
114+
@Override
115+
public void reset() {
116+
farIdMap.clear();
117+
reverseFarIdMap.clear();
118+
bufferFarIds.clear();
119+
farIdToUeAddrs.clear();
120+
nextGlobalFarId = 0;
121+
}
122+
123+
@Override
124+
public Map<UpfRuleIdentifier, Integer> getFarIdMap() {
125+
return Map.copyOf(farIdMap.asJavaMap());
126+
}
127+
128+
@Override
129+
public int globalFarIdOf(UpfRuleIdentifier farIdPair) {
130+
int globalFarId = farIdMap.compute(farIdPair,
131+
(k, existingId) -> {
132+
return Objects.requireNonNullElseGet(existingId, () -> nextGlobalFarId++);
133+
}).value();
134+
log.info("{} translated to GlobalFarId={}", farIdPair, globalFarId);
135+
return globalFarId;
136+
}
137+
138+
@Override
139+
public int globalFarIdOf(ImmutableByteSequence pfcpSessionId, int sessionLocalFarId) {
140+
UpfRuleIdentifier farId = new UpfRuleIdentifier(pfcpSessionId, sessionLocalFarId);
141+
return globalFarIdOf(farId);
142+
143+
}
144+
145+
@Override
146+
public String queueIdOf(int schedulingPriority) {
147+
return (SCHEDULING_PRIORITY_MAP.get(schedulingPriority)).toString();
148+
}
149+
150+
@Override
151+
public String schedulingPriorityOf(int queueId) {
152+
return (SCHEDULING_PRIORITY_MAP.inverse().get(queueId)).toString();
153+
}
154+
155+
@Override
156+
public UpfRuleIdentifier localFarIdOf(int globalFarId) {
157+
return reverseFarIdMap.get(globalFarId);
158+
}
159+
160+
public void learnFarIdToUeAddrs(PacketDetectionRule pdr) {
161+
UpfRuleIdentifier ruleId = UpfRuleIdentifier.of(pdr.sessionId(), pdr.farId());
162+
farIdToUeAddrs.compute(ruleId, (k, set) -> {
163+
if (set == null) {
164+
set = new HashSet<>();
165+
}
166+
set.add(pdr.ueAddress());
167+
return set;
168+
});
169+
}
170+
171+
@Override
172+
public boolean isFarIdBuffering(UpfRuleIdentifier farId) {
173+
checkNotNull(farId);
174+
return bufferFarIds.contains(farId);
175+
}
176+
177+
@Override
178+
public void learBufferingFarId(UpfRuleIdentifier farId) {
179+
checkNotNull(farId);
180+
bufferFarIds.add(farId);
181+
}
182+
183+
@Override
184+
public void forgetBufferingFarId(UpfRuleIdentifier farId) {
185+
checkNotNull(farId);
186+
bufferFarIds.remove(farId);
187+
}
188+
189+
@Override
190+
public void forgetUeAddr(Ip4Address ueAddr) {
191+
farIdToUeAddrs.keySet().forEach(
192+
farId -> farIdToUeAddrs.computeIfPresent(farId, (farIdz, ueAddrs) -> {
193+
ueAddrs.remove(ueAddr);
194+
return ueAddrs;
195+
}));
196+
}
197+
198+
@Override
199+
public Set<Ip4Address> ueAddrsOfFarId(UpfRuleIdentifier farId) {
200+
return farIdToUeAddrs.getOrDefault(farId, Set.of()).value();
201+
}
202+
203+
@Override
204+
public Set<UpfRuleIdentifier> getBufferFarIds() {
205+
return Set.copyOf(bufferFarIds);
206+
}
207+
208+
@Override
209+
public Map<UpfRuleIdentifier, Set<Ip4Address>> getFarIdToUeAddrs() {
210+
return Map.copyOf(farIdToUeAddrs.asJavaMap());
211+
}
212+
213+
// NOTE: FarIdMapListener is run on the same thread intentionally in order to ensure that
214+
// reverseFarIdMap update always finishes right after farIdMap is updated
215+
private class FarIdMapListener implements MapEventListener<UpfRuleIdentifier, Integer> {
216+
@Override
217+
public void event(MapEvent<UpfRuleIdentifier, Integer> event) {
218+
switch (event.type()) {
219+
case INSERT:
220+
reverseFarIdMap.put(event.newValue().value(), event.key());
221+
break;
222+
case UPDATE:
223+
reverseFarIdMap.remove(event.oldValue().value());
224+
reverseFarIdMap.put(event.newValue().value(), event.key());
225+
break;
226+
case REMOVE:
227+
reverseFarIdMap.remove(event.oldValue().value());
228+
break;
229+
default:
230+
break;
231+
}
232+
}
233+
}
234+
}

0 commit comments

Comments
 (0)