webFilter实现mock接口

简介: 这段代码实现了一个名为 `MockFilter` 的类,继承自 `WebFilter` 接口,用于处理 HTTP 请求和响应。它通过从 Redis 缓存中获取配置信息来决定是否使用模拟数据或缓存数据来响应请求。如果开启了生产模式或关闭了模拟和缓存功能,则直接放行请求。否则,它会检查请求体并根据配置返回相应的模拟或缓存数据。同时,该过滤器支持对响应结果进行处理,并将结果存储回 Redis 中。

package com.ph.sp.gateway.filter;

import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.lang.Pair;
import cn.hutool.core.util.StrUtil;
import lombok.extern.slf4j.Slf4j;
import org.reactivestreams.Publisher;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.core.annotation.Order;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferFactory;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import org.springframework.data.redis.core.HashOperations;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.http.HttpHeaders;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.http.server.reactive.ServerHttpRequestDecorator;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.http.server.reactive.ServerHttpResponseDecorator;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ServerWebExchange;
import org.springframework.web.server.WebFilter;
import org.springframework.web.server.WebFilterChain;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import javax.annotation.Resource;
import java.nio.charset.StandardCharsets;
import java.util.Map;
import java.util.concurrent.atomic.AtomicReference;

@Component
@Order(3)
@Slf4j
@SuppressWarnings("all")
public class MockFilter implements WebFilter {

@Value("${mockClose:true}")
private boolean mockClose;
@Value("${cacheClose:true}")
private boolean cacheClose;
@Value("${isPrd:true}")
private boolean prd;
private final String cache = "A_CACHE_";
private final String mock = "A_MOCK_";
@Resource
private RedisTemplate<String, Map<String, String>> hashTemplate;

@Override
public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) {
    if (prd || mockClose && cacheClose) {
        return chain.filter(exchange);
    }
    ServerHttpRequest request = exchange.getRequest();
    String uri = request.getURI().getPath();
    Map<String, String> mockConfig = getRedisConfig(exchange, mock);
    Map<String, String> cacheConfig = getRedisConfig(exchange, cache);
    if (CollUtil.isEmpty(mockConfig) && CollUtil.isEmpty(cacheConfig)) {
        return chain.filter(exchange);
    }
    AtomicReference<String> requestBodyContent = new AtomicReference<>("");
    Flux<DataBuffer> body = exchange.getRequest().getBody();
    return body.doOnNext(buffer -> {
        byte[] bytes = new byte[buffer.readableByteCount()];
        buffer.read(bytes);
        DataBufferUtils.release(buffer);
        requestBodyContent.set(new String(bytes, StandardCharsets.UTF_8));
    }).then(Mono.defer(() -> diyFilter(requestBodyContent.get(), exchange, chain, mockConfig, cacheConfig)
    )).then();
}

private Mono<Void> diyFilter(String body, ServerWebExchange exchange, WebFilterChain chain,
                             Map<String, String> mockConfig, Map<String, String> cacheConfig) {
    String uri = exchange.getRequest().getURI().getPath();
    ServerHttpResponse response = exchange.getResponse();
    Pair<String, String> mockResult = check(mockConfig, body);
    if (StrUtil.isNotBlank(mockResult.getValue())) {
        return hit(mock, mockResult.getKey(), Base64Utils.decode(mockResult.getValue()), uri, response);
    }
    Pair<String, String> cacheResult = check(cacheConfig, body);
    if (StrUtil.isNotBlank(cacheResult.getValue())) {
        return hit(cache, cacheResult.getKey(), cacheResult.getValue(), uri, response);
    }
    Flux<DataBuffer> cachedFlux = Flux.defer(() -> Mono.just(exchange.getResponse().bufferFactory().wrap(body.getBytes())));
    ServerHttpRequest mutatedRequest = new ServerHttpRequestDecorator(exchange.getRequest()) {
        @Override
        public HttpHeaders getHeaders() {
            HttpHeaders headers = new HttpHeaders();
            headers.putAll(exchange.getRequest().getHeaders());
            headers.remove(HttpHeaders.CONTENT_LENGTH);
            headers.setContentLength(body.getBytes().length);
            return headers;
        }

        @Override
        public Flux<DataBuffer> getBody() {
            return cachedFlux;
        }
    };
    if (StrUtil.isNotBlank(cacheResult.getKey())) {
        ServerHttpResponseDecorator mutatedResponse = decoratedResponse(exchange, cacheResult.getKey());
        return chain.filter(exchange.mutate().request(mutatedRequest).response(mutatedResponse).build());
    }
    return chain.filter(exchange.mutate().request(mutatedRequest).build());
}

private Mono<Void> hit(String type, String key, String val, String uri, ServerHttpResponse resp) {
    log.info("命中缓存数据:(key: {}, hashKey:{},用完请手动删除该条hash,谢谢)", type, type + uri, key);
    resp.getHeaders().setContentType(HeaderConstants.APPLICATION_JSON_UTF8);
    return resp.writeWith(Mono.fromSupplier(() -> resp.bufferFactory().wrap(val.getBytes(StandardCharsets.UTF_8))));
}

private ServerHttpResponseDecorator decoratedResponse(ServerWebExchange exchange, String key) {
    String path = exchange.getRequest().getURI().getPath();
    ServerHttpResponse originalResponse = exchange.getResponse();
    DataBufferFactory bufferFactory = originalResponse.bufferFactory();
    return new ServerHttpResponseDecorator(originalResponse) {
        @Override
        public Mono<Void> writeWith(Publisher<? extends DataBuffer> body) {
            if (body instanceof Mono) {
                Mono<? extends DataBuffer> mono = (Mono<? extends DataBuffer>) body;
                body = mono.flux();
            }
            if (body instanceof Flux) {
                Flux<? extends DataBuffer> fluxBody = (Flux<? extends DataBuffer>) body;
                return super.writeWith(fluxBody.buffer().map(dataBuffer -> {
                    DataBufferFactory dataBufferFactory = new DefaultDataBufferFactory();
                    DataBuffer join = dataBufferFactory.join(dataBuffer);
                    byte[] content = new byte[join.readableByteCount()];
                    join.read(content);
                    DataBufferUtils.release(join);
                    HashOperations<String, String, String> operation = hashTemplate.opsForHash();
                    operation.put(cache + path, key, new String(content, StandardCharsets.UTF_8));
                    originalResponse.getHeaders().setContentLength(content.length);
                    return bufferFactory.wrap(content);
                }));
            }
            return super.writeWith(body);
        }
    };
}

private Map<String, String> getRedisConfig(ServerWebExchange exchange, String pre) {
    String uri = exchange.getRequest().getURI().getPath();
    HashOperations<String, String, String> operation = hashTemplate.opsForHash();
    return operation.entries(pre + uri);
}

private Pair<String, String> check(Map<String, String> config, String body) {
    return config.entrySet().stream().filter(e -> StrUtil.contains(body, e.getKey())).findFirst()
            .map(e -> Pair.of(e.getKey(), e.getValue())).orElse(Pair.of(null, null));
}

}

相关文章
|
9天前
|
人工智能 JSON API
全网刷屏的 Jev 模型正式开放!一手实战测评 + 保姆级教程
全网爆火的 Jev 模型是什么?有什么用?怎么使用?怎么接入 AI 编程工具?效果真的好么?傻子可懂的 Jev 保姆级实战教程 + 项目实战测评来啦
7714 13
|
8天前
|
人工智能 测试技术 API
最近全网爆火的 Jev 到底是什么?适合干什么、怎么用,一篇讲透!
Jev是TypeSafe AI推出的“系统一模型”,不生成文本,专做毫秒级结构化决策:Choice(多选)、Score(打分)、Noul(是非概率)。响应快193倍、成本低444倍,适合工单路由、内容审核、测试定级等高频判断场景。
1661 4
最近全网爆火的 Jev 到底是什么?适合干什么、怎么用,一篇讲透!
|
5天前
|
人工智能 JavaScript 芯片
DeepSeek 官方偷偷上传 Harness 桌面端安装包,我已经用上了。。附最新下载地址
DeepSeek Harness 官方的桌面端安装包被网友扒出来了,2 分钟讲明白如何使用,体验如何,适合作为 AI 编程工具么?附最新 Windows 和 Mac 双端的下载地址
1434 1
|
8天前
|
人工智能 并行计算 PyTorch
秋叶 ComfyUI 2026 整合包 v3.2 完整部署教程:Python 3.13 + Torch 2.13 全栈升级
秋叶aaaki ComfyUI 2026年8月整合包v3.2正式发布!全面升级Python 3.13.11、PyTorch 2.13.0+cu130及ComfyUI v0.30.2,原生支持MiniMax H3、Wan 2.2、Qwen-Image-2.1等2026主流音视频/图像模型,解压即用,无需环境配置。
1243 10
|
21天前
|
人工智能 自然语言处理 安全
阿里云千问办公 QwenWork详细介绍:产品核心能力、典型场景、价格及常见问题解答
千问办公是阿里云推出的一站式AI办公平台,主打"不止于对话,更注重交付",依托通义千问旗舰大模型,用户一句话即可完成数据分析、PPT生成、视频剪辑等复杂任务,直接输出可用成果。产品深度打通钉钉生态与企业OA,覆盖桌面端、网页端,提供企业标准版198元/人/月等多档订阅方案,新用户注册即赠2000积分,适配工程师、HR、财务等多职业办公场景,成为能动手干活的"全能AI同事"。
3686 10
|
6天前
|
人工智能 编解码 并行计算
MiniMax-H3 一键整合包技术文档:8G 显存运行 AI 漫剧制作 —— 角色替换 / 动作迁移 / 文图生视频部署与调参指南
MiniMax H3 是 MiniMax 开源的全模态视频生成模型,支持文/图/音/视多条件输入,输出最高2K、15秒带双声道音频视频。本文档详述其Int8量化版在8GB显存下的本地一键部署、三段式工作流(EDIT/REPLACE/CONTINUE)、参数调优及常见问题排查。(239字)
|
16天前
|
缓存 IDE Java
【保姆级】Android Studio下载、安装和汉化教程(2026最新)
Android Studio 是 Google 官方推出的免费 Android 应用开发集成环境,基于 IntelliJ IDEA,内置模拟器、调试器、性能分析及 Compose 界面工具,功能全面,文档丰富,是安卓开发首选工具。(239字)
1759 1

热门文章

最新文章