前言:

  在上一篇,我们的项目从简单的producer-consumer模型成功转化为了一个拥有模拟网关,通过NAT+PORT配合虚拟机实现拟真的跨端通讯的较为贴切现实生产的形态,同时利用动态配置以及长连接改造拥有了一定的多权限并发的模拟能力。但是仅仅满足于Windows和Linux是明显背离项目全链路模拟初衷的,而且producer本身并没有完成和Windows网关的隔离,也没能发送除了模拟信号外的有效信息——拟真的目标并没有达到可用水平,况且现实中的client-server模型,client和server本身也应当承担多种身份而非单一的发送/接受,那么也需要我们进行相应的改造。

  本篇我们将通过引入Android端client,改造flask网关与Linux上的consumer来解决上述问题,并引入一个验证码-用户权限业务场景,让项目达到一个基本用的状态。

一.完整链路的搭建

  Android+Windows+Linux的平台组合在同一台Windows宿主机上面同时运行带来的调试和环境搭建等问题较多,同时当前利用的NAT+PORT模式,跨平台转发会涉及到一些比较麻烦的问题,所以我会先在这里说一下,但要注意,具体情况具体分析,一些数据可能会因为环境的不同改变 。

  该文章基于雷电模拟器(LD)AndroidStudio(AS)作为Android-client的开发和运行环境

  假设你拥有了一个基本的使用OKHTTP模块编写的Android网络程序(当然如果没有我下文也会给出一个简单的示例),那么下一步就应当是考虑如何将目前孤立的Android-client和Windows上的flask网关进行有效链接,我们先回顾一下之前用pika建立的链接

......
credentials = pika.PlainCondentials(USER,PASSWORD)
parameter = pika.ConnectionParameters(HOST_ADDR,credentials = credentials)
connection = pika.BlockingConnection(parameter)
channel = connection.channel()
......
channel.close()

  但是在离开现成的库后,Android和Win上的flask网关之间建立网络联系就需要更复杂一点——由于NAT+PORT模式,flask网关获取信息的功能实际上是通过监听对应port来实现的,而这里的port本身属于Windows宿主机,同时LD并不是和VM一样的专业虚拟机平台,建立链接不能简单通过localaddress实现——reverse反向映射是个比较可行的方法。我这里运用了AS作为IDE自带的ADB功能,以免和LD的ADB功能冲突,下面我会结合具体powershell指令来演示并结合注释讲解:

cd D:Android\Sdk\platform-tools  ##SDK的安装路径

.\adb.exe devices -l ##显示和ADB端口connected的设备

(当然,如果你找不到adb,可以先试试where.exe adb)
运行后如果显示为空表,当然,第一次建立conn前大概率是空的,就可以参考下方程序

cd C:\leidian\LDPlayer9  ##运行实例的目录
.\ldconsole.exe list2  ##显示各实例状态

根据你的具体情况决定,我这里提供一份模版

Get-NetTCPConnection -State Listen |
Where-Object { $_.LocalPort -ge 5555 -and $_.LocalPort -le 5600 } |
Sort-Object LocalPort |
Format-Table LocalAddress, LocalPort, OwningProcess  ##查找ADB端口

cd D:Android\Sdk\platform-tools  ##重定位

.\adb.exe connect 127.0.0.1:5555  ##让LD实例链接到5555端口
.\adb.exe devices -l ##验证
.\adb.exe -s 127.0.0.1:5555 reverse tcp:5000 tcp:5000  ##正式建立反向映射
.\adb.exe -s 127.0.0.1:5555 recerse --list  ##检测端口上的映射设备

具体运行示例如图,这样一个基本的反向映射就建立完成了。当然,后续开多个实例的流程也类似。值得注意的一点是,区分不同port的类型很重要,flask网关监听的port和LD实例中的虚拟port,反代后的revers port容易混淆,而且会随你的个人设置而和示例有出入。

  Windows和Linux的链接仍旧使用之前文章的思路,借助RabbitMQ实现queue通讯,这里不过多赘述,如果有疑问,可以去翻阅之前文章,那里有更详细的思路和相关flask示例。

二.Linux-server与Android-client的双向沟通

  在之前的文章中我提到过一个隐性问题:当前producer-consumer模型中,组件之间的身份是单一的——Linux上的consumer只能接受queue的message并操作MySQL进行一些简单操作,Windows上的producer只能发送message——模型关系虽然完整,但是现实生产环境中这个模式过于简化了。在C-S模式里,client和server,包含中间的网关gateway,往往会同时作为consumer和producer,根据需求而改变身份,比如一个简单的联网APP,client在发送request后也需要从producer转化为consumer处理接受到的response,同样的server也需要类似的身份转化。因此我们必须改造现有模式,这样才能让项目有可接受的拟真效果。

1.一个简单的Android-client

  TIP:如果你刚刚开始接触ASIDE和Kotlin编程,那么大概率你会被AS的gradle配置折磨,如果是,我建议你修改对应的配置,先把gradle的zip给下到本地,然后通过修改project的gradle的路径,直接使用本地URL,当然你也可以配置mirror网站,只不过有一些已经失效了。同时如果你对Kotlin感到陌生,我建议你去先学Java或者干脆边练边学,我本人之前就是边练边学的,而且Kotlin本身教程比较少,如果想找网课,建议学Java,它们的关系可以类比C++/Python,技多不压身。

  在具体代码示例前我先讲一些基本知识:在AS中开发Android应用相比于我们之前的pycharm需要多加一些操作:之前的flask网关直接在terminal中操作即可,但是client不行,至少在我们当前的环境中,AS本身是不负责用它的VirtualDevice测试的,这种方法对个人电脑负担比较大而且需要额外了解JVM等更复杂的知识,效果也不是很理想,而LD的测试需要我们把写好的代码打包为APK,然后让LD安装,这样才能让它接入之前配置好的LD环境。下面我会结合具体代码演示client的基本逻辑:

package com.example.conmunication_Test  ##project名称

import android.app.Activity  ##引入相关的库和控件,如果你用过QT我觉得你会感到有些熟悉
import android.os.Bundle
import android.widget.Button
import android.widget.EditText
import android.widget.TextView
import okhttp3.Call
import okhttp3.Callback
import okhttp3.MediaType.Companion.toMediaType
import okhttp3.OkHttpClient
import okhttp3.Request
import okhttp3.RequestBody.Companion.toRequestBody
import okhttp3.Response
import org.json.JSONObject
import java.io.IOException

class MainActivity : Activity() {

    private val client = OkHttpClient()

    private lateinit var messageInput: EditText  ##declare空间中具体的类
    private lateinit var sendButton: Button      ##lateinit var可以让这些类不需要
    private lateinit var resultText: TextView    ##在declare时赋值,区别于var
                                                 
    private var activeCall: Call? = null  ##将call信号和null关联

    override fun onCreate(savedInstanceState: Bundle?) {  ##理解为创建一个UI实例
        super.onCreate(savedInstanceState)                ##让input的数据和APP后台互通
        setContentView(R.layout.activity_main)

        messageInput = findViewById(R.id.messageInput)  ##通过id查找对应的View
        sendButton = findViewById(R.id.sendButton)
        resultText = findViewById(R.id.resultText)

        sendButton.setOnClickListener {
            sendMessage()
        }
    }

    private fun sendMessage() {  ##发送message
        val message = messageInput.text.toString().trim()  ##读取UI中input的message

        if (message.isEmpty()) {
            resultText.text = "Message cannot be empty"  ##非空判定
            return
        }

        activeCall?.cancel()  ##改变call的status
        sendButton.isEnabled = false  
        resultText.text = "Requesting..."

        val json = JSONObject()  ##JSON格式的message模式
            .put("message", message)
            .toString()

        val requestBody = json.toRequestBody(  ##赋值message给requestbody
            "application/json; charset=utf-8".toMediaType()
        )

        val request = Request.Builder()  ##开始初始化request的具体逻辑和各种变量
            .url("http://127.0.0.1:5000/message")
            .post(requestBody)
            .build()

        val call = client.newCall(request)  ##call开始
        activeCall = call

        call.enqueue(object : Callback {  

            override fun onFailure(call: Call, error: IOException) { ##错误处理
                finishRequest(  
                    call,  
                    "Request failed: ${error.message}"
                )
            }

            override fun onResponse(call: Call, response: Response) {  
                val statusCode = response.code  ##response_status显示
                val body = response.use {
                    it.body?.string().orEmpty()  ##非空判定
                }

                val displayText = formatResponse(  ##显示回传的response内容
                    statusCode,
                    body,
                )

                finishRequest(call, displayText)  
            }
        })
    }

    private fun formatResponse(  ##组装reponse
        statusCode: Int,
        body: String,
    ): String {
        if (body.isBlank()) {  ##非空判定
            return "HTTP $statusCode\n\nEmpty response"
        }

        return try {  ##开始接受并组装text
            val json = JSONObject(body)

            val code = json.optString("code", "")
            val status = json.optString("status", "")
            val reason = json.optString("reason", "")
            val found = json.optBoolean("found", false)

            when {
                code.isNotBlank() -> {
                    "Verification code: $code"  ##status检验,code!=null
                }

                status == "ok" && !found -> {
                    "Confirm code not found\n\n$body"
                }

                status.isNotBlank() -> {  ##debug信息,可以类比logcat
                    "Status: $status" +
                        if (reason.isNotBlank()) {
                            "\nReason: $reason"
                        } else {
                            ""
                        } +
                        "\n\n$body"
                }

                else -> {  
                    "HTTP $statusCode\n\n$body"  ##内容组装=true
                }
            }
        } catch (error: Exception) {  ##debug依据
            "HTTP $statusCode\n\n$body"
        }
    }

    private fun finishRequest(  ##request的结束函数,方便减少工程量和debug
        call: Call,
        message: String,
    ) {
        runOnUiThread {
            if (activeCall !== call) {
                return@runOnUiThread  ##lambda用法,让单个子功能不影响全局
            }                         ##可以结合协程,channel等kotlin知识了解 
                                      ##这里@是为了只退出内层函数
            activeCall = null  ##重置activeCall为null,结束Call

            if (isFinishing || isDestroyed) {
                return@runOnUiThread
            }

            sendButton.isEnabled = true
            resultText.text = message
        }
    }

    override fun onDestroy() {  ##彻底结束request流程
        activeCall?.cancel()
        activeCall = null
        super.onDestroy()
    }
}

   提一嘴@,这里是kotlin的lambda用法,防止直接return外层example函数,如果你接触过Java类语言,注意不要和另一个@的用法,注解,给混淆了,同时这里省略了XML布局,但是XML比较简单,我这里不再给出。

2.关于Windows-flask-gateway

  之前的gateway可以保留大量的可复用部分,Linux和Windows之间的通讯场景在之前的文章中已经得到了有效验证和简单测压,目前欠缺的就是回传到Android-client的部分,所以我现在给出的代码只是一部分,其它部分可以和前文略加改造后通用,并且由于下文会对flask进一步改造,所以当前代码可以当作了解逻辑和整体构造的素材:

    ......##基础配置和import

def publish_test_message(message_text: str):
    message_id = str(uuid.uuid4())  ##message_id生成

    envelope = {  ##信封构造
        "message_id": message_id,
        "source": "android",
        "type": TYPE_TEST,
        "created_at": datetime.now(timezone.utc).isoformat(),  ##code有效期检验
        "payload": {
            "message": message_text,
        },
    }

    body = json.dumps(  ##从envelope获得json格式的body
        envelope,
        ensure_ascii=False,
    ).encode("utf-8")  ##注意编码一致性

    ......##pika调用,可以参考Linux和Windows-gateway的链接建立

def request_confirm_code(confirm_code_id: str):  ##验证码获取
    request_id = str(uuid.uuid4())  ##验证码生成

    ......##pika调用,可以参考Linux和Windows-gateway的链接建立

            envelope = {  ##信封构造
            "message_id": request_id,
            "source": "android",
            "type": TYPE_CONFIRM_CODE_REQUEST,
            "created_at": datetime.now(timezone.utc).isoformat(),
            "payload": {
                "confirm_code_id": confirm_code_id,
            },
        }

        body = json.dumps(  ##内容拆解
            envelope,
            ensure_ascii=False,
        ).encode("utf-8")

        ......##和Linux-server的链接,可以理解为consumer模块,参考之前的Linux-consumer
       
    finally:
        if connection is not None and connection.is_open:
            connection.close()

@app.get("/health")  ##status检测
def health():
    return jsonify(status="ok")

@app.post("/message")  ##message推送
def message():
    data = request.get_json(silent=True)  

    if not isinstance(data, dict):  ##格式检验,当前为dict(字典)
        return jsonify(
            status="error",
            reason="invalid_json",
        ), 400

    message_text = data.get("message")  ##获取message并赋值给text

    if not isinstance(message_text, str) or not message_text.strip():  
        return jsonify(  ##非空检测
            status="error",
            reason="message_required",
        ), 400

    message_text = message_text.strip()

    if message_text.startswith("confirm_code:"):  ##验证码校验
        confirm_code_id = message_text.removeprefix(
            "confirm_code:"
        ).strip()

        if not confirm_code_id:  ##验证码非空检测
            return jsonify(
                status="error",
                reason="confirm_code_id_required",
            ), 400

        try:
            result = request_confirm_code(confirm_code_id)  ##结果调用
        except pika.exceptions.AMQPError as error:  ##debug用日志
            app.logger.exception("RabbitMQ request failed")
            return jsonify(
                status="error",
                reason="rabbitmq_error",
                detail=str(error),
            ), 503

        if result is None:  ##过期检测
            return jsonify(
                status="timeout",
                reason="consumer_response_timeout",
            ), 504

        return jsonify(result), 200

    try:  
        message_id = publish_test_message(message_text)
    except pika.exceptions.AMQPError as error:  ##检测是否为rabbitmq问题
        app.logger.exception("RabbitMQ publish failed")
        return jsonify(
            status="error",
            reason="rabbitmq_error",
            detail=str(error),
        ), 503

    return jsonify(  ##成功获取
        accepted=True,
        message_id=message_id,
    ), 202


if __name__ == "__main__":  
    app.run(  ##通过之前的reverse配置,port反向映射,让message被LD实例获取,传给client
        host="127.0.0.1",
        port=5000,
        use_reloader=False,
    )

3.Linux-server改造和可能的问题

  和上方的Windows-flask-gateway类似,server也不需要大规模改造,具体改造可以参考基本的Linux-Windows之间的flask+RabbitMQ链接,所以我在这里不会给出具体代码,但我仍会提醒两个可能犯错的地方:

VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)
ON DUPLICATE KEY UPDATE
    message_id = VALUES(message_id)

  如你所见,在这部分的SQL语句中有9个需要的赋值的变量,也就是not null,由于我们之前大量使用了自己declare的json格式,比如message,envelope等,如果你因为某些原因,比如变量名错误,以外静默等,导致parameter不满足占位要求,那么就会在Linux的MySQL/consumer那里出现大量报错刷屏,也就是message在RabbitMQ的queue中为durable,一直作为ready状态导致大量尝试并报错——建议先把message从管理页面踢出,然后慢慢改,不然会很烦

  然后是关于utf-8的,老实说这个报错很隐秘:如果你发现LD上的APP点击后就不断闪退,那么可能就是栈溢出问题,而栈溢出问题就很可能是utf-8导致的——在encode/编译时,比如下方:

 private val jsonMediaType =
        "application/json; charset=utf-8".toMediaType()

  如果你不小心把charset=utf-8哪怕用空格分开,那么就会导致链式反应,结果就是APP闪退,让之前的各种debug措施毫无办法——后续文章我会专门讲这种隐性报错

三.多实例多权限情况下通讯骨干的优化

  当前的开发中,运行多个实例会带来一个问题:权限的配置应当如何完成?一个比较好的方法是复用之前的confirm_code,运用Linux-server作一次”权限申请“。同时,当前的checker,也就是Windows-flask-gateway还是只是packet检验,也就是检测包的合法性和对应的level来做静态的”死“判定,这在现实中的用户权限动态改变的环境下是低效的。下面我会分别给出对应的回答。

1.权限动态配置

  通过confirm_code,我们可以联合Linux-server的处理功能形成一个简单的流程:新增一个client的upgrade功能,通过request的message中带有的level申请需求,让Linux-server识别关键的role类别,通过role的种类来发送一个特定形式的验证码——比如level = 1代表role = normal_uesr,confirm_code = 1xxxxx——client接受到code后通过将code输入到新的envelope变量中,让该信封的level通过code被网关识别并处理,通过和client的新增变量结合——每一个LD实例通过APK标识一个producer_id,这样就能够实现一个有效的动态配置。虽然这种临时协议并不适用生产环境,但对于简单的实验已经足够了。

  当然,client应当带有一个默认的配置,满足格式,比如统一默认为normal_user,之后再通过主动的request进行upgrade,这样可以减少工程量,也便于维护——多个实例可以使用同一个APK,分发和debug也可以锁定这个唯一的APK,也让后续更多功能的实现不必频繁面临APK之间细微的不同导致的debug困难。下面我会给出简化的功能代码,你可以根据它改造上面的基础client:

   ......
private fun requestPermissionUpgrade() {  ##申请权限升级
        val requestedRole = requestedRoleInput  ##申请的目标
            .text
            .toString()
            .trim()

        if (requestedRole.isEmpty()) {  ##非空检测
            resultText.text = "requested role cannot be empty"
            return
        }

        val body = JSONObject()  ##内容JSON构造
            .put("type", TYPE_PERMISSION_REQUEST)
            .put("producer_id", producerId)
            .put("session_id", sessionId)
            .put("requested_role", requestedRole)

        postJson("/message", body) { statusCode, responseBody ->  ##发送函数
            if (statusCode !in 200..299) {  ##状态码校验
                resultText.text =
                    "HTTP $statusCode\n\n$responseBody"
                return@postJson  ##lamda
            }

            try {  
                val json = JSONObject(responseBody)  ##level_code获取
                val upgradeCode = json.optString(
                    "request_level_code",
                    "",
                )

                if (upgradeCode.isBlank()) {  ##非空检验
                    resultText.text =
                        "Upgrade code missing\n\n$responseBody"
                    return@postJson
                }

                upgradeCodeInput.setText(upgradeCode)  ##获得的权限码

                resultText.text = buildString {  ##结果内容构造
                    append("Upgrade code received: ")
                    append(upgradeCode)
                    append("\n\n")
                    append(responseBody)
                }
            } catch (error: Exception) {  ##debug日志依据
                resultText.text =
                    "Invalid response:\n\n$responseBody"
            }
        }
    }

    ......
private fun postJson(    ##发送函数,参考之前的基础client
        path: String,
        body: JSONObject,
        onResult: (Int, String) -> Unit,
    ) {
        activeCall?.cancel()
        setBusy(true)

        val requestBody = body
            .toString()
            .toRequestBody(jsonMediaType)
                                            ##一些配置的declare,参考之前的基础client
        val request = Request.Builder()
            .url(BASE_URL + path)
            .post(requestBody)
            .build()

        val call = client.newCall(request)
        activeCall = call

        call.enqueue(object : Callback {

            override fun onFailure(
                call: Call,
                error: IOException,
            ) {
                finishRequest(
                    call,
                    "Request failed: ${error.message}",
                )
            }

            override fun onResponse(
                call: Call,
                response: Response,
            ) {
                val statusCode = response.code
                val responseBody = response.use {
                    it.body?.string().orEmpty()
                }

                runOnUiThread {
                    if (activeCall !== call) {
                        return@runOnUiThread
                    }

                    activeCall = null

                    if (isFinishing || isDestroyed) {
                        return@runOnUiThread
                    }

                    setBusy(false)  ##状态控制
                    onResult(statusCode, responseBody)
                }
            }
        })
    }
    ......
 private companion object {  ##默认配置区
        const val BASE_URL = "http://127.0.0.1:5000"

        const val PRODUCER_ID_KEY = "producer_id"
        const val TOKEN_KEY = "access_token"  ##token限制,后续我会结合网关讲
        const val ROLE_KEY = "role"

        const val DEFAULT_ROLE = "normal_user"

        const val TYPE_MESSAGE = "message"
        const val TYPE_PERMISSION_REQUEST =
            "permission_request"
        const val TYPE_PERMISSION_ACTIVATE =
            "permission_activate"
    }

  2.网关的改造

  如上所述,Windows-flask-gateway应当承担动态配置功能,具体实施起来其实比较简单,甚至免除了level的计算改用更简单的检验比较,只不过需要我们注意讲client的envelope中相应的关键变量role给准确传递到相应位置。我们在这里探讨的重点在于检验模式的扩展——token。

  如果你接触过AI,那么你肯定多少知道一些token的知识,但我在这里讲的是一些不要花钱的东西,和AI关系不是很大,下面我们来看一下token的工作流程,这里它的中文名”令牌“或许更符合它在这里的功能:

normal_user
->申请升级
->client验证申请
->返回一次性升级凭证
->client提交升级凭证
->server签发 access_token
->后续request携带 access_token
Android 携带 token
->Flask 验证 token
->Flask 得到 principal_id、role、scope
->Flask 将身份信息写入内部
->RabbitMQ
->Consumer (linux-server)

  通过用户在上文提到的权限申请功能,server生成一个主体身份principal_id,然后签发token发送给client,client通过token绑定自己的身份从而获得相应的权限,之后request会带有token,这样就可以利用token来更方便高效的实现更复杂的检验,比如事务,时间,版本号等,下面我会结合一些代码,帮助你了解更具体的运行逻辑,方便你自己进行改造

  

def activate_permission(data: dict):  ##激活权限
    producer_id = data.get("producer_id")
    session_id = data.get("session_id")  ##获取request的要求和身份
    upgrade_code = data.get("request_level_code")

    if not isinstance(producer_id, str) or not producer_id.strip():  ##非空检测
        return jsonify(
            status="error",
            reason="producer_id_required",
        ), 400

    if not isinstance(session_id, str) or not session_id.strip():
        return jsonify(
            status="error",
            reason="session_id_required",
        ), 400

    if not isinstance(upgrade_code, str) or not upgrade_code.strip():
        return jsonify(
            status="error",
            reason="request_level_code_required",
        ), 400

    upgrade_code = upgrade_code.strip()  ##获取权限码

    with identity_lock:
        upgrade = pending_upgrades.get(upgrade_code)  

        if upgrade is None:  ##非空/过期校验
            return jsonify(
                status="error",
                reason="invalid_or_used_code",
            ), 400

        if upgrade["expires_at"] <= time.time():  ##升级检测
            pending_upgrades.pop(upgrade_code, None)
            return jsonify(
                status="error",
                reason="upgrade_code_expired",
            ), 400

        if upgrade["producer_id"] != producer_id:  ##身份检测
            return jsonify(
                status="error",
                reason="upgrade_code_owner_mismatch",
            ), 403

        pending_upgrades.pop(upgrade_code, None)  ##组装

        role = upgrade["requested_role"]  ##token生成
        access_token = secrets.token_urlsafe(32)
        expires_at = time.time() + ACCESS_TOKEN_TTL_SECONDS

        access_grants[access_token] = {  ##token内容
            "producer_id": producer_id,
            "session_id": session_id,
            "role": role,
            "expires_at": expires_at,
        }

    app.logger.info(  ##日志
        "permission activated: producer_id=%s role=%s",
        producer_id,
        role,
    )

    return jsonify(
        status="ok",
        role=role,
        access_token=access_token,
        expires_in=ACCESS_TOKEN_TTL_SECONDS,
    ), 200


def resolve_role(  ##规划role,方便行为决策
    producer_id: str,
    access_token: str,
):
    if not access_token:  ##无token就返回default的配置
        return ROLE_NORMAL

    now = time.time()

    with identity_lock:  ##本体检测
        grant = access_grants.get(access_token)

        if grant is None:
            return None

        if grant["expires_at"] <= now:  ##过期校验
            access_grants.pop(access_token, None)
            return None

        if grant["producer_id"] != producer_id:
            return None

        return grant["role"]

  当然,Linux-server对应的consumer也需要结合token模式进行对应的改造,逻辑都是类似的,你可有通过添加对于token的接受和检验,并结合MySQL做出更具体/个性化的协议

  当前单access_token仅仅是做测试,生产环境往往需要更复杂的防伪,储存,校验,同时token的种类也多种多样,如果需要可以自由扩展,后期我也会结合MySQL做相应的专项优化。

尾记:

  从项目正式开工到现在大概是三周,项目从零变成了一个基本能够自由扩展的框架,可以做到基本的全链路:多实例,并发,权限分级,网关校验......功能开始变得更贴合现实,但是还不够。后续我计划深入优化并添加更多功能,让拟真程度更上一层楼——下一步估计会是MySQL数据库的主场,顺便做一个完整的压力测试——当然,如果你看到这里,记得少喝咖啡,喝多了感觉更困了。

附上第一份运行成功的截图

Logo

一站式 AI 云服务平台

更多推荐