内核和用户进程协作 — 同步阻塞

用户进程—>内核:

用户进程发出创建socket指令—>切换到内核态创建socket并初始化

内核—>用户进程:

网络包到达网卡,通过硬中断和软中断处理—>唤醒用户进程

同步阻塞IO总体流程:

用户进程发出IO请求后,内核会查看struct socksk_receive_queue中有无数据:

  • 没有数据,内核就将用户线程加入到struct socket 的等待队列,将用户进程标记为阻塞状态(TASK_INTERRUPTTIBLE),通过调度器交出CPU,进入睡眠状态
  • 数据就绪后,调用sk—>sk_data_ready唤醒等待队列中的用户进程,进程状态变为TASK_RUNNING,被调度器重新加入运行队列


在这里插入图片描述

1.1 socket相关内容

socket内核结构:

在这里插入图片描述

用户态通过socket函数创建一个socket,内核在内部会通过系统调用SYSCALL_DEFINE3(socket, int, family, int, type, int, protocol)相应创建一系列的socket相关的内核对象

创建socket的主要函数__sock_create

// net/socket.c
int __sock_create(struct net *net, int family, int type, int protocol,
			 struct socket **res, int kern)
{
	int err;
	struct socket *sock;
	const struct net_proto_family *pf;
	...
    // 分配socket对象
	sock = sock_alloc();
	...
    // 获得每个协议族的操作表  
	pf = rcu_dereference(net_families[family]);
	...
    // 调用指定协议族的创建函数,对于AF_INET对应的是inet_create
	err = pf->create(net, sock, protocol, kern);
	...
}

对于AF_INET协议对应的创建函数是inet_create:根据类型SOCK_STREAM查找到对于TCP定义的操作方法实现集合inet_stream_ops和tcp_prot,并把它们分别设置到socket->ops和sock->sk_prot上

// net/ipv4/af_inet.c
static struct inet_protosw inetsw_array[] =
{
	{
		.type =       SOCK_STREAM,
		.protocol =   IPPROTO_TCP,
		.prot =       &tcp_prot,
		.ops =        &inet_stream_ops,
		.no_check =   0,
		.flags =      INET_PROTOSW_PERMANENT |
			      INET_PROTOSW_ICSK,
	},
	...
};

static int inet_create(struct net *net, struct socket *sock, int protocol,
		       int kern)
{
	struct sock *sk;
	struct inet_protosw *answer;
	struct inet_sock *inet;
	struct proto *answer_prot;  
	...
	list_for_each_entry_rcu(answer, &inetsw[sock->type], list) {

		err = 0;
		if (protocol == answer->protocol) {
			if (protocol != IPPROTO_IP)
				break;
		} else {
			if (IPPROTO_IP == protocol) {
				protocol = answer->protocol;
				break;
			}
			if (IPPROTO_IP == answer->protocol)
				break;
		}
		err = -EPROTONOSUPPORT;
	}
	...
    // 将inet_stream_ops赋到socket->ops上
	sock->ops = answer->ops;
    // 获取tcp_prot
	answer_prot = answer->prot;
	...
    // 分配sock对象,并把tcp_prot赋到sock->sk_prot上
	sk = sk_alloc(net, PF_INET, GFP_KERNEL, answer_prot);
	...
	// 对sock对象进行初始化
	sock_init_data(sock, sk);
	...
}

初始化sock对象:

当软中断上收到数据包时通过调用sk_data_ready函数指针来唤醒在socket上等待的进程

// net/core/sock.c
void sock_init_data(struct socket *sock, struct sock *sk)
{
	...
	sk->sk_data_ready	=	sock_def_readable;    // 设置数据就绪回调函数,当有数据可读时调用
	sk->sk_write_space	=	sock_def_write_space; // 设置写空间可用回调函数,当发送缓冲区有空间时调用
	sk->sk_error_report	=	sock_def_error_report;// 设置错误报告回调函数,当发生错误时调用
	...
}

1.2 用户进程等待接收消息

通过系统调用,用户进程进入内核态,执行一系列的内核协议层函数,接着在socket对象的接收队列中查看是否有数据,无数据则将自己添加到socket等待队列中,让出CPU。

在这里插入图片描述

1.2.1 代码分析

执行recvfrom(用于从套接字接收数据)系统调用:

// net/socket.c
SYSCALL_DEFINE6(recvfrom, int, fd, void __user *, ubuf, size_t, size,
		unsigned int, flags, struct sockaddr __user *, addr,
		int __user *, addr_len)
{
	struct socket *sock;
	...
    // 根据用户传入的fd找到socket对象
	sock = sockfd_lookup_light(fd, &err, &fput_needed);
	...
	err = sock_recvmsg(sock, &msg, size, flags); //从指定的 struct socket 中接收数据,调用具体协议的接收操作
	...
}

sock_recvmsg—>__sock_recvmsg —> _ _sock_recvmsg_nosec:调用socket对象ops中的recvmsg,recvmsg指向的是inet_recvmsg

// net/socket.c
static inline int __sock_recvmsg_nosec(struct kiocb *iocb, struct socket *sock,
				       struct msghdr *msg, size_t size, int flags)
{
	...
	return sock->ops->recvmsg(iocb, sock, msg, size, flags);
}

调用socket对象里sk_prot下的recvmsg方法,recvmsg方法对应的是tcp_recvmsg方法

// net/ipv4/af_inet.c
int inet_recvmsg(struct kiocb *iocb, struct socket *sock, struct msghdr *msg,
		 size_t size, int flags)
{
	...
	err = sk->sk_prot->recvmsg(iocb, sk, msg, size, flags & MSG_DONTWAIT,
				   flags & ~MSG_DONTWAIT, &addr_len);
	...
}

接收队列为空或者接收的数据不够多,则调用sk_wait_data将此进程阻塞

// net/ipv4/tcp.c
int tcp_recvmsg(struct kiocb *iocb, struct sock *sk, struct msghdr *msg,
		size_t len, int nonblock, int flags, int *addr_len)
{
	...
	int copied = 0;
	...
	do {
		...
		// 遍历接收队列接收数据
		skb_queue_walk(&sk->sk_receive_queue, skb) {
			...
                
		}
		...
		if (copied >= target) {
			release_sock(sk);
			lock_sock(sk);
		} else // 没有收到足够数据,启用sk_wait_data阻塞当前进程
			sk_wait_data(sk, &timeo);
		...
	} while (len > 0);
	...
}

用户进程如何将自己阻塞sk_wait_data
在这里插入图片描述

// net/core/sock.c
int sk_wait_data(struct sock *sk, long *timeo)
{
	int rc;
	// 当前进程(current)关联到所定义的等待队列项上
	DEFINE_WAIT(wait);

 	// 调用sk_sleep获取sock对象下的wait
    // 并准备挂起,将当前进程设置为可打断(INTERRUPTIBLE)
	prepare_to_wait(sk_sleep(sk), &wait, TASK_INTERRUPTIBLE);
	set_bit(SOCK_ASYNC_WAITDATA, &sk->sk_socket->flags);
	// 通过调用schedule_timeout让出CPU,然后进行睡眠
	rc = sk_wait_event(sk, timeo, !skb_queue_empty(&sk->sk_receive_queue));
	...
}

DEFINE_WAIT宏定义:定义一个等待队列项wait

// include/linux/wait.h
#define DEFINE_WAIT_FUNC(name, function)				
	wait_queue_t name = {			//定义一个wait_queue_t类型的变量name,表示一个等待队列项			
		.private	= current,		//存储当前进程的struct task_struct	
		.func		= function,		//回调函数
		.task_list	= LIST_HEAD_INIT((name).task_list),	//将 wait_queue_t 链接到等待队列头(wait_queue_head_t)的链表中
	}

#define DEFINE_WAIT(name) DEFINE_WAIT_FUNC(name, autoremove_wake_function)

调用sk_sleep获取socket对象下的等待队列列表头wait_queue_head_t

// include/net/sock.h
static inline wait_queue_head_t *sk_sleep(struct sock *sk)
{
	BUILD_BUG_ON(offsetof(struct socket_wq, wait) != 0);
	return &rcu_dereference_raw(sk->sk_wq)->wait;
}

调用prepare_to_wait来把新定义的等待队列项wait插入sock对象的等待队列

// kernel/wait.c
void prepare_to_wait(wait_queue_head_t *q, wait_queue_t *wait, int state)
{
	unsigned long flags;

	wait->flags &= ~WQ_FLAG_EXCLUSIVE;
	spin_lock_irqsave(&q->lock, flags);
	if (list_empty(&wait->task_list))
		__add_wait_queue(q, wait);
	set_current_state(state);
	spin_unlock_irqrestore(&q->lock, flags);
}

1.3 软中断模块

软中断接收数据的过程:

​ 软中断中收到数据包后,若是TCP数据包则对应执行tcp_v4_rcv函数;若是ESTABLISH状态下的数据包,则最终会把数据拆出来放到对应socket的接收队列中,然后调用sk_data_ready来唤醒用户进程。
在这里插入图片描述

1.3.1 代码分析

根据网络包中的header中的source和dest信息在本机上查询对应的socket:

// net/ipv4/tcp_ipv4.c
int tcp_v4_rcv(struct sk_buff *skb)
{
	...
	// 获取tcp header
	th = tcp_hdr(skb);
	// 获取ip header
	iph = ip_hdr(skb);
	...
	// 根据数据包header中的IP、端口信息查找到对应的socket
	sk = __inet_lookup_skb(&tcp_hashinfo, skb, th->source, th->dest);
	...
	// socket未被用户锁定
	if (!sock_owned_by_user(sk)) {
		...
		{
			if (!tcp_prequeue(sk, skb))
				ret = tcp_v4_do_rcv(sk, skb);
		}
	}
	...
}

调用tcp_v4_do_rcv函数:

// net/ipv4/tcp_ipv4.c
int tcp_v4_do_rcv(struct sock *sk, struct sk_buff *skb)
{
	...
	if (sk->sk_state == TCP_ESTABLISHED) {
		...
		// 执行连接状态下的数据处理
		if (tcp_rcv_established(sk, skb, tcp_hdr(skb), skb->len)) {
			rsk = sk;
			goto reset;
		}
		return 0;
	}
	// 其他非ESTABLISH状态的数据包处理
	...
}
// net/ipv4/tcp_input.c
int tcp_rcv_established(struct sock *sk, struct sk_buff *skb,
			const struct tcphdr *th, unsigned int len)
{
				...
				// 接收数据放到队列中
				eaten = tcp_queue_rcv(sk, skb, tcp_header_len,
						      &fragstolen);
			...
			// 数据准备好,唤醒socket上阻塞掉的进程
			sk->sk_data_ready(sk, 0);
			...
}
// net/ipv4/tcp_input.c
static int __must_check tcp_queue_rcv(struct sock *sk, struct sk_buff *skb, int hdrlen,
		  bool *fragstolen)
{
	...
    // 把接收到的数据放到socket的接收队列的尾部  
	if (!eaten) {
		__skb_queue_tail(&sk->sk_receive_queue, skb);
		skb_set_owner_r(skb, sk);
	}
	return eaten;
}

sock_init_data函数已经把sk_data_ready指针设置成了sock_def_readable函数,默认就是就绪处理函数

// net/core/sock.c
static void sock_def_readable(struct sock *sk, int len)
{
	struct socket_wq *wq;

	rcu_read_lock();
	wq = rcu_dereference(sk->sk_wq);
    // 有进程在此socket的等待队列
	if (wq_has_sleeper(wq))
        // 唤醒等待队列上的进程
		wake_up_interruptible_sync_poll(&wq->wait, POLLIN | POLLPRI |
						POLLRDNORM | POLLRDBAND);
	sk_wake_async(sk, SOCK_WAKE_WAITD, POLL_IN);
	rcu_read_unlock();
}

唤醒在socket上因为等待数据而被阻塞掉的进程:

// include/linux/wait.h
#define wake_up_interruptible_sync_poll(x, m)				
	__wake_up_sync_key((x), TASK_INTERRUPTIBLE, 1, (void *) (m))
// kernel/sched/core.c
void __wake_up_sync_key(wait_queue_head_t *q, unsigned int mode,
			int nr_exclusive, void *key)
{
	unsigned long flags;
	int wake_flags = WF_SYNC;

	if (unlikely(!q))
		return;

	if (unlikely(!nr_exclusive))
		wake_flags = 0;

	spin_lock_irqsave(&q->lock, flags);
	__wake_up_common(q, mode, nr_exclusive, wake_flags, key);
	spin_unlock_irqrestore(&q->lock, flags);
}

__wake_up_common实现唤醒。函数调用的参数nr_exclusion传入1,只唤醒一个进程。

__wake_up_common中找出一个等待队列项curr,然后调用其curr->func。在上一部分(等待接收消息)中,使用DEFINE_WAIT()定义等待队列项时,内核把curr->func设置成了autoremove_wake_function

// kernel/sched/core.c
static void __wake_up_common(wait_queue_head_t *q, unsigned int mode,
			int nr_exclusive, int wake_flags, void *key)
{
	wait_queue_t *curr, *next;

	list_for_each_entry_safe(curr, next, &q->task_list, task_list) {
		unsigned flags = curr->flags;

		if (curr->func(curr, mode, wake_flags, key) &&
				(flags & WQ_FLAG_EXCLUSIVE) && !--nr_exclusive)
			break;
	}
}

autoremove_wake_function中,调用了default_wake_function

// kernel/sched/core.c
int default_wake_function(wait_queue_t *curr, unsigned mode, int wake_flags,
			  void *key)
{
	return try_to_wake_up(curr->private, mode, wake_flags);
}

调用try_to_wake_up,这个函数执行完的时候,在socket上等待而被阻塞的进程就被推入可运行队列里了,这又将产生一次进程上下文切换的开销

1.3 同步阻塞IO方式缺点:

  • 上下文切换开销大:一个进程专门为等待一个socket上的数据而被从CPU上拿下来,换上另一个进程;数据准备好后,睡眠的进程又会被唤醒;总共产生两次上下文切换开销
  • 模型单一低效:socket和进程一对一
Logo

「智能机器人开发者大赛」官方平台,致力于为开发者和参赛选手提供赛事技术指导、行业标准解读及团队实战案例解析;聚焦智能机器人开发全栈技术闭环,助力开发者攻克技术瓶颈,促进软硬件集成、场景应用及商业化落地的深度研讨。 加入智能机器人开发者社区iRobot Developer,与全球极客并肩突破技术边界,定义机器人开发的未来范式!

更多推荐