Modify logic for header copy from input to output
This is primarily related to Cloud Events. Since we delegate to a separate class for post processing, if outpt message is Cloud Event we will not be doing anything to with regard to header copy in SimpleFunctionRegistry and unstead delegate it to CloudEventFunctionInvocationHelper
This commit is contained in:
@@ -1,150 +0,0 @@
|
||||
/*
|
||||
* Copyright 2020-2020 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* https://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package io.spring.cloudevent;
|
||||
|
||||
import java.net.URI;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
|
||||
import org.springframework.boot.SpringApplication;
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
import org.springframework.boot.web.client.RestTemplateBuilder;
|
||||
import org.springframework.cloud.function.cloudevent.CloudEventHeaderEnricher;
|
||||
import org.springframework.cloud.function.cloudevent.CloudEventMessageBuilder;
|
||||
import org.springframework.cloud.function.cloudevent.CloudEventMessageUtils;
|
||||
import org.springframework.cloud.function.web.util.HeaderUtils;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.http.RequestEntity;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.MessageBuilder;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* Sample application that demonstrates how user functions can be triggered by cloud event.
|
||||
* Events can come from anywhere (e.g., HTTP, Messaging, RSocket etc).
|
||||
* Given that this particular sample comes already with spring-cloud-function-web support each
|
||||
* function is a valid REST endpoint where function name signifies URL path (e.g., http://localhost:8080/asPOJOMessage).
|
||||
*
|
||||
* Simply start the application and post cloud event to individual function - (see individual 'curl' command at each function).
|
||||
*
|
||||
* You can also run CloudeventDemoApplicationTests.
|
||||
*
|
||||
* @author Oleg Zhurakousky
|
||||
*
|
||||
*/
|
||||
@SpringBootApplication
|
||||
public class CloudeventDemoApplication {
|
||||
|
||||
boolean consumerSuccess;
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
SpringApplication.run(CloudeventDemoApplication.class, args);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<Message<String>, String> asStringMessage() {
|
||||
return v -> {
|
||||
System.out.println("Received Cloud Event with raw data: " + v);
|
||||
return v.getPayload();
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@Bean
|
||||
public Function<String, String> asString() {
|
||||
return v -> {
|
||||
System.out.println("Received raw Cloud Event data: " + v);
|
||||
return v;
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@Bean
|
||||
public Function<Message<SpringReleaseEvent>, String> asPOJOMessage() {
|
||||
return v -> {
|
||||
System.out.println("Received Cloud Event with POJO data: " + v);
|
||||
return v.getPayload().toString();
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@Bean
|
||||
public Function<SpringReleaseEvent, String> asPOJO() {
|
||||
return v -> {
|
||||
System.out.println("Received POJO Cloud Event data: " + v);
|
||||
return v.toString();
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<Message<SpringReleaseEvent>, Message<SpringReleaseEvent>> consumeAndProduceCloudEvent() {
|
||||
return ceMessage -> {
|
||||
SpringReleaseEvent data = ceMessage.getPayload();
|
||||
data.setVersion("2.0");
|
||||
data.setReleaseDateAsString("01-10-2006");
|
||||
|
||||
return MessageBuilder.withPayload(data).build();
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public CloudEventHeaderEnricher cloudEventHeaderEnricher() {
|
||||
return headers -> {
|
||||
return headers.setSource("https://interface21.com/").setType("com.interface21");
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@Bean
|
||||
public Function<Map<String, Object>, Map<String, Object>> consumeAndProduceCloudEventAsMapToMap() {
|
||||
return ceMessage -> {
|
||||
ceMessage.put("version", "10.0");
|
||||
ceMessage.put("releaseDate", "01-10-2050");
|
||||
return ceMessage;
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<SpringReleaseEvent, SpringReleaseEvent> consumeAndProduceCloudEventAsPojoToPojo() {
|
||||
return event -> {
|
||||
event.setVersion("2.0");
|
||||
return event;
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Consumer<Message<SpringReleaseEvent>> pojoConsumer(CloudEventHeaderEnricher enricher, RestTemplateBuilder builder) {
|
||||
return eventMessage -> {
|
||||
Message<?> newMessage = enricher.enrich(CloudEventMessageBuilder.fromMessage(eventMessage)).build(CloudEventMessageUtils.DEFAULT_ATTR_PREFIX);
|
||||
RequestEntity<SpringReleaseEvent> entity = RequestEntity.post(URI.create("http://foo.com"))
|
||||
.headers(HeaderUtils.fromMessage(newMessage.getHeaders()))
|
||||
.body(eventMessage.getPayload());
|
||||
List<String> sourceHeader = entity.getHeaders().get("ce-source");
|
||||
Assert.isTrue(sourceHeader.get(0).equals("https://interface21.com/"), "'source' must be https://interface21.com/");
|
||||
List<String> typeHeader = entity.getHeaders().get("ce-type");
|
||||
Assert.isTrue(typeHeader.get(0).equals("com.interface21"), "'source' must be com.interface21");
|
||||
List<String> idHeader = entity.getHeaders().get("ce-id");
|
||||
Assert.notEmpty(idHeader, "'id' must not be null");
|
||||
List<String> specversionHeader = entity.getHeaders().get("ce-specversion");
|
||||
Assert.notEmpty(specversionHeader, "'specversion' must not be null");
|
||||
this.consumerSuccess = true;
|
||||
};
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,77 +0,0 @@
|
||||
/*
|
||||
* Copyright 2020-2020 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* https://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package io.spring.cloudevent;
|
||||
|
||||
import java.text.ParseException;
|
||||
import java.text.SimpleDateFormat;
|
||||
import java.util.Date;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonFormat;
|
||||
|
||||
/**
|
||||
* An example POJO that represents cloud event data
|
||||
*
|
||||
* @author Oleg Zhurakousky
|
||||
*
|
||||
*/
|
||||
public class SpringReleaseEvent {
|
||||
|
||||
@JsonFormat(shape = JsonFormat.Shape.STRING, pattern = "dd-MM-yyyy")
|
||||
private Date releaseDate;
|
||||
|
||||
private String releaseName;
|
||||
|
||||
private String version;
|
||||
|
||||
public Date getReleaseDate() {
|
||||
return releaseDate;
|
||||
}
|
||||
|
||||
public void setReleaseDate(Date releaseDate) {
|
||||
this.releaseDate = releaseDate;
|
||||
}
|
||||
|
||||
public void setReleaseDateAsString(String releaseDate) {
|
||||
try {
|
||||
this.releaseDate = new SimpleDateFormat("dd-MM-yyyy").parse(releaseDate);
|
||||
}
|
||||
catch (ParseException e) {
|
||||
throw new IllegalArgumentException(e);
|
||||
}
|
||||
}
|
||||
|
||||
public String getReleaseName() {
|
||||
return releaseName;
|
||||
}
|
||||
|
||||
public void setReleaseName(String releaseName) {
|
||||
this.releaseName = releaseName;
|
||||
}
|
||||
|
||||
public String getVersion() {
|
||||
return version;
|
||||
}
|
||||
|
||||
public void setVersion(String version) {
|
||||
this.version = version;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "releaseDate:" + new SimpleDateFormat("dd-MM-yyyy").format(releaseDate) + "; releaseName:" + releaseName + "; version:" + version;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user