7 C multi wątki i procesy

PROCESY WĄTKI SYNCHRONIZACJA

Procesy w linuxie

# top    - lista procesów

# ps      - procesy skojarzone z uzytkownikiem

# ps  -aux    - wszystkie procesy (a-dla wszystkich uzytkowników,u-szczegółowo,x-demony)



Otrzymanie swojego ID i ID rodzica w C:
#include <stdio.h>
#include <stdlib.h> // EXIT_SUCCESS
#include <unistd.h> // getpid(), getppid()

int main(int argc, char* argv[])
{
        printf("pid = %d\n", getpid());
        printf("ppid = %d\n", getppid());

        exit(EXIT_SUCCESS);
}


Do uruchomienia nowego procesu z wnętrza programu (np inny plik uruchamialny) służy EXECL
po zakonczeniu nie wraca do źródłowego programu:
printf("moj_prog = %d, uruchamiamy inny program\n", getpid());

        //no PATH, full path needed
        execl("/bin/ps", "ps", "-u", NULL);

Do uruchamiania z powrotem służy polecenie system
int main(int argc, char* argv[])
{
        printf("moj_prog = %d, uruchamiamy inny program\n", getpid());
        system("ps -u");
        printf("moj_prog = %d, konczymy prace\n", getpid());

        exit(EXIT_SUCCESS);
}


FORK
Klonuje nasz program na dwa procesy i zwraca:
-1   błąd
0    w kodzie dziecka
0> w kodzie rodzica (tak naprawde rodzic otrzymuje nr ID dziecka)


Tworzymy zmienna np "pid"
pid_t pid;
i uruchamiamy na niej fork
pid=fork();
ewentualnie w jednej linii jak w przykładzie

#include <stdio.h>
#include <stdlib.h> // EXIT_SUCCESS
#include <unistd.h> // getpid(), getppid()

int main(int argc, char* argv[])
{
        pid_t pid = fork();
        if(pid == 0)
        {
                printf("child pid = %d\n",getpid());
                printf("child parent pid = %d\n", getppid());
                exit(EXIT_SUCCESS);
        }
        else if(pid > 0)
        {
                printf("pid ktory dostal rodzic(dziecka?) = %d\n",pid);
                printf("parent pid = %d\n",getpid());
                printf("parent parent pid = %d\n",getppid());
                exit(EXIT_SUCCESS);
        }
        else
        {
                printf("error");
        }
}

efekt uruchomienia:
(pierwszy uruchamia sie rodzic)

zmienna! pid ktory dostal rodzic(dziecka) = 4032
parent pid = 4031                          (prawdziwe id rodzica)
parent parent pid = 3192                   (id basha)

user$ child pid = 4032 (prawdziwe id dziecka)
child parent pid = 1448

waitpid
blokuje wykonanie procesu do czasu zakonczenia poddanego, synchronizuje

#include <stdio.h>
#include <stdlib.h>   // EXIT_SUCCESS
#include <unistd.h>   // getpid(), getppid()
#include <sys/wait.h> // wait(), waitpid()

int main(int argc, char* argv[])
{
        int pid = fork();
        int status;

        if(pid == -1)
        {
                printf("error\n");
        }
        else if(pid == 0)
        {
                printf("child %d\n", getpid());
                sleep(10);
                printf("child %d | konczymy prace\n", getpid());
                exit(EXIT_SUCCESS);
        }
        else
        {
                printf("parent %d | czekamy na zakonczenie child\n", getpid());
                waitpid(pid, &status, 0);
                printf("parent %d | po wait\n", getpid());
                exit(EXIT_SUCCESS);
        }
}

efekt uruchomienia:

parent 4079 | czekamy na zakonczenie child
child 4080
child 4080 | konczymy prace
parent 4079 | po wait

Polecenie "wait" zwraca ID dziecka na które czekało  " ret=wait(&status); "
Przykład stworzenia dziesięciu dzieci i czekanie na nie:

#include <stdio.h>
#include <stdlib.h>   // EXIT_SUCCESS
#include <unistd.h>   // getpid(), getppid()
#include <sys/wait.h> // wait(), waitpid()

int main(int argc, char* argv[])
{
        int status, ret;

        for(int i = 0 ; i < 10 ; i++)
        {
                int pid = fork();

                if(pid == 0)
                {
                        printf("child %d\n", getpid());
                        sleep(10+i);
                        printf("child %d - konczymy prace\n", getpid());
                        exit(EXIT_SUCCESS);
                }
        }

        sleep(1);
        for(int i = 0 ; i < 10 ; i++)
        {
                printf("parent - czekamy na zakonczenie child\n");
                ret = wait(&status);
                printf("parent - po wait ret = %d | status = %d\n", ret, status);
        }

        printf("parent - w koncu koniec\n");
        exit(EXIT_SUCCESS);
}

I jako wynik mamy 10 nr child a potem sekunda po sekundzie kolejne 10 zakończeń (pętla dziecka dodaje po 1 do każdego następnego)

child 4177
child 4169
child 4168
parent - czekamy na zakonczenie child
child 4168 - konczymy prace
parent - po wait ret = 4168 | status = 0
parent - czekamy na zakonczenie child
child 4169 - konczymy prace
parent - po wait ret = 4169 | status = 0
parent - czekamy na zakonczenie child
child 4170 - konczymy prace
parent - po wait ret = 4170 | status = 0

Zombi - gdy rodzic działa, dziecko zakonczyło, ale on nie wie:
(np dziecko sleep(1);  a rodzic czyta  getchar(); przez kilka sekund
na nowym terminalu z poleceniem # ps  -aux | grep Z
widzimy proces zombi czyli rodzica z zamknietym dzieckiem

Sierota - proces dziecka któremu zakończył pracę rodzić (dziecko dłużej sleep niż rodzić) wtedy dziecko jak zostanie sierota i zapytamy o getppid to odpowie ID powłoki ). Tak się tworzy Deamony które działają w tle pod władzą roota.

Wątki

Część programu wykonywana współbieżnie w obrębie jednego procesu.

wątki zwane są „procesami lekkimi” (ang. light weight processes)
   • w jednym procesie może istnieć wiele wątków
   • w ramach jednego procesu wątki współdzielą przestrzeń adresową oraz struktury
   systemowe (Wątek współdzieli pamięć, ale ma osobny tylko stos, proces ma osobną pamięć)


WAŻNE
podczas kompilacji trzeba dodać:   -Lpthread 


Wyświetlanie wprocesu wraz z wątkami w Linuxie:
# ps -eLf  | grep  a.out


Tworzenie wątku polecenie   " pthread_create (4 arg) "
jeśli zwraca 0 - OK

#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>  // EXIT_SUCCESS
#include <pthread.h> // pthread

void* thread_fun(void *arg)
{
        printf("New thread id = %li\n", pthread_self());
        pthread_exit(NULL);
}


int main(int argc, char* argv[])
{
        pthread_t thread_id;
        int res;

        res = pthread_create(&thread_id,NULL,&thread_fun,NULL);
/*
          &thread_id,  // pthread_t *thread / id nowego watku
          NULL,          // atrybuty watku, NULL - domyslne
          &thread_fun, //  funkcja w ktorej watek bedzie dzialal
          NULL         // void *arg/ argument przekazany do funkcji w ktorej watek bedzie dzialal
                                                
*/
        if(0 == res) {
                printf("Main thread id = %li | new thread = %li\n", pthread_self(), thread_id);
        } else {
                printf("Thread creation error\n");
        }

        sleep(5);
        exit(EXIT_SUCCESS);
}

Mamy też polecenie " pthread_self() "  czyli pobranie własny id wątku:
program najpierw wyświetla własny + ID nowego
czeka 5 sek
i uruchamiana jest Funkcja nowostworzonego wątku (widać ten sam ID)

Main thread id = 140554533377792 | new thread = 140554525067008
New thread id = 140554525067008

Inne polecenia:
pthread_exit  poprawne zakończenie wątku


pthread_join czekanie aż wątek się zakończy, zwraca 0 gdy wskazany wątek zakonczy swoje działanie, przykład:

#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>  // EXIT_SUCCESS
#include <pthread.h> // pthread

void* thread_fun(void *arg)
{
        printf("New thread\n");
        sleep(5);
        printf("New thread ends\n");
        pthread_exit(NULL);
}

int main(int argc, char* argv[])
{
        pthread_t thread_id;

        if(0 == pthread_create(&thread_id, NULL, &thread_fun, NULL))
        {
                printf("Main threads waits...\n");
                if(0 == pthread_join(thread_id, NULL))
                        printf("Main threads ends\n");
                else
                        printf("pthread_join error\n");
        }
        else
        {
                printf("thread_create error\n");
        }

        exit(EXIT_SUCCESS);
}



Kolejny przykład, wysyłanie argumentu do funkcji wywołanej w wątku(4 parametr)
Przykład bez warunków chroniących przed błędami podczas tworzenia.

#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>  // EXIT_SUCCESS
#include <pthread.h> // pthread

void* thread_fun(void *arg)
{
        printf("thread arg 1 = %d\n", *(int*)(arg));
        *(int*)(arg) = 10;
        printf("thread arg 1 = %d\n", *(int*)(arg));
        pthread_exit(NULL);
}

int main(int argc, char* argv[])
{
        pthread_t thread_id;
        int temp = 5;

        pthread_create(&thread_id, NULL, thread_fun, &temp);
        pthread_join(thread_id, NULL);

        printf("main thread temp = %d\n", temp);
        exit(EXIT_SUCCESS);
}

Wynik działania, jak widać zmiana zmiennej wewnątrz funkcji wpłynęła na zmienna poza funkcją (głównego wątku)

thread arg 1 = 5
thread arg 1 = 10
main thread temp = 10


Kolejny przykład jak wątki mogą zwracać
fun_1: przez zmianę zmiennych w strukturze
fun_2: przez zwrot parametru zwykłym return
fun_3 przez zwrot parametru jako argument pthread_exit(argument)

#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>  // EXIT_SUCCESS
#include <pthread.h> // pthread

typedef struct Sum {
        int a;
        int b;
        int result;
} Sum;

void* fun_1(void* arg)
{
        Sum* p = (Sum*)arg;
        p->result = p->a + p->b;
        pthread_exit(NULL);
}

void* fun_2(void* arg)
{
        Sum* p = (Sum*)arg;
        p->result = p->a + p->b;
        return (p);
}

void* fun_3(void* arg)
{
        Sum* p = (Sum*)arg;
        p->result = p->a + p->b;
        pthread_exit(p);
}

int main(int argc, char* argv[])
{
        pthread_t thread1, thread2, thread3;
        Sum first  = {2,3};
        Sum second = {2,3};
        Sum third  = {2,3};
        Sum* status1;
        Sum* status2;

        pthread_create(&thread1, NULL, fun_1, &first);
        pthread_create(&thread2, NULL, fun_2, &second);
        pthread_create(&thread3, NULL, fun_3, &third);

        pthread_join(thread1, NULL);
        pthread_join(thread2, (void**)&status1); // void **retval
        pthread_join(thread3, (void**)&status2); // void **retval

        printf("thread 1 first  : %d + %d = %d\n", first.a, first.b, first.result);
        printf("thread 2 second : %d + %d = %d\n", second.a, second.b, status1->result);
        printf("thread 3 third  : %d + %d = %d\n", third.a, third.b, status2->result);          

        exit(EXIT_SUCCESS);
}

Wynik działania:

thread 1 first  : 2 + 3 = 5
thread 2 second : 2 + 3 = 5
thread 3 third  : 2 + 3 = 5

Przykład programu otwierający 1000 wątków.

#include<stdio.h>
#include<pthread.h>
#include<unistd.h>  // sleep


void *wyswietl(void *arg)
{
  while(1)
  {
    printf("POSIX C wyswietlam swoj pid= %li\n",pthread_self());
//    sleep(1);
  }
    //   pthread_exit(NULL);
}

int main()
{


pthread_t  myThreads[1000];                          
int i=0;
while(i<1000)
{    


pthread_create(&myThreads[i],NULL,&wyswietl,NULL);              

//sleep(1);
i++;
}

int b;                    
scanf("%d",&b);       //nie konczy maina           

return 0;
}





Synchronizacja i komunikacja

Podstawowe pojęcia
•Komunikacja(ang. Communication) –mechanizmy służące do wymiany danch pomiędzy procesami/wątkami

•Synchronizacja(ang. Synchronization) –mechanizmy służące do zgrania działania procesów/wątków

•Sygnały(ang. Signals) –mechanizmy programowych przerwań mogące być wykorzystwane w roli komunikacji lub synchronizacji w roli komunikacji lub synchronizacji


Uproszone zalecenie

Wątki:
• mutexy
• zmienne warunkowe

Procesy powiązane:
• semafory nienazwane

Procesy niepowiązane:
• semafory nazwane



Mutex

Pojęcia:
Operacja atomowa (Atomic op.)- operacja, która na określonym poziomie abstrakcji jest niepodzielna.

Instrukcje mikroprocesora z punktu widzenia systemu operacyjnego są operacjami atomowymi - gdy procesor rozpocznie wykonywanie instrukcji nie można jej przerwać ani w żaden sposób wpłynąć na jej realizację.

Sekcja krytyczna (Critical section) – fragment kodu programu, w którym korzysta się z zasobu dzielonego, a co za tym idzie w danej chwili może być wykorzystywany przez co najwyżej jeden wątek

Wyścig (Race condition) - sytuacja gdzie współdzielone dane są odczytywane jednocześnie przez różne procesy i rezultat czasu dostępu zależy od ich prędkości działania.

Wzajemne wykluczenie (Mutual exclusion) - mechanizm chroniący sekcję krytyczną. Kiedy jeden proces/wątek jest w sekcji krytycznej to mechanizm blokuje dostęp innym procesom/ wątkom.

Zagłodzenie - proces nigdy nie wybierany przez shedulera
Zakleszczenie - Deadlock - oba procesy czekają jeden na drugiego
Zakleszczenie - Livelock - oba procesy nieustannie zmieniają swój stan pod wpływem drugiego

Mutex mechanizm służący do synchronizacji działa wątków poprzez atomowy dostęp do współdzielonego zasobu (ang. shared resource)

• wzajemne wykluczanie (ang. mutual exclusion)

dwa stany:
• zablokowany (ang. locked) – ponowne blokowanie powoduje zablokowania działania wątku
• odblokowany (ang. unlocked) odblokować mutex może tylko ten sam wątek, co go blokował




pthread_mutex_init(3p)          inicjalizacja mutexu
pthread_mutex_lock(3p)        blokowanie mutexu
pthread_mutex_unlock(3p)    odblokowanie mutexu
pthread_mutex_destroy(3p)   zniszczenie mutexu



Generalny przykład wywoływania i działania mutexu

#include<stdlib.h>
#include<pthread.h>
#include<unistd.h>  // sleep

pthread_mutex_t mutex;   //zmienna mutexu


void funkcja1(void)     //funkcja ktora na watku bedzie blok. mutex
{
pthread_mutex_lock( &mutex );   //blokuje mutex
//petla
pthread_mutex_unlock( &mutex );  //odblokowywuje mutex
}


int main ()
{

pthread_mutex_init(&mutex, NULL);    //inicjalizacja mutexu

pthread_t ping;                     //tworzenie watku
pthread_create( &ping, NULL, (void*)&funkcja1, NULL );   //uruchamianie funkcji w watku


pthread_join( ping, NULL );    //czeka az watek pink zakonczy
pthread_mutex_destroy(&mutex);  //niszczy mutex

return 0;
}




Przykład bez synchronizacji!!


#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>  // EXIT_SUCCESS
#include <pthread.h> // pthread
#include <string.h>  // strcpy(), strcpy()

#define BUF_SIZE 100
char buff[BUF_SIZE];

void* fun_write(void *arg)
{
 int i = 0;
 while(1) {
  if(i % 2) {
   memcpy(buff, "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", BUF_SIZE);
  } else {
   memcpy(buff, "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", BUF_SIZE);
  }
  i++;
 }
 pthread_exit(NULL);
}

void* fun_read(void* arg)
{
 char buff_local[BUF_SIZE];

 while(1) {
  memcpy(buff_local, buff, BUF_SIZE); 
  printf("%s\n", buff_local);
 }
 pthread_exit(NULL);
}

int main(int argc, char* argv[])
{
 pthread_t thread_write, thread_read;

 pthread_create(&thread_write, NULL, fun_write, NULL);
 pthread_create(&thread_read,  NULL, fun_read,  NULL);

 pthread_join(thread_write, NULL);
 pthread_join(thread_read,  NULL);
 
 exit(EXIT_SUCCESS);
}


//Przykład z Mutex ********************************************************************

#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>  // EXIT_SUCCESS
#include <pthread.h> // pthread
#include <string.h>  // strcpy(), strcpy()

#define BUF_SIZE 100
char buff[BUF_SIZE];
pthread_mutex_t my_mutex;

void* fun_write(void *arg)
{
 int i = 0;
 while(1) {
  pthread_mutex_lock(&my_mutex);
   if(i % 2) {
    memcpy(buff, "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", BUF_SIZE);
   } else {
    memcpy(buff, "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", BUF_SIZE);
   }
  pthread_mutex_unlock(&my_mutex);
  i++;
 }
 pthread_exit(NULL);
}
void* fun_read(void* arg)
{
 char buff_local[BUF_SIZE];

 while(1) {
  pthread_mutex_lock(&my_mutex);
   memcpy(buff_local, buff, BUF_SIZE); 
   printf("%s\n", buff_local);
  pthread_mutex_unlock(&my_mutex);
 }
 pthread_exit(NULL);
}

int main(int argc, char* argv[])
{
 pthread_t thread_write, thread_read;
 
 // inicjalizacja mutexu
 pthread_mutex_init(&my_mutex, NULL);

 pthread_create(&thread_write, NULL, fun_write, NULL);
 pthread_create(&thread_read,  NULL, fun_read,  NULL);
 
 pthread_join(thread_write, NULL);
 pthread_join(thread_read,  NULL);
 
 // zniszczenie mutexu
 pthread_mutex_destroy(&my_mutex);
 
 exit(EXIT_SUCCESS);
}
przykłady M1/003thread_sync.../ 002

przykład pinpong


#include <pthread.h>
#include <stdio.h>
#include <string.h>

void ping_function(void);
void pong_function(void);

int flag=0;
pthread_mutex_t mutex;

void ping_function(void)
{
  do {
    pthread_mutex_lock( &mutex );
    if( flag )
      {
        printf("ping\n");
        flag = 0;
      }
    pthread_mutex_unlock( &mutex );
  } while(1);

  pthread_exit(0);
 }


void pong_function(void)
{
  do {
    pthread_mutex_lock( &mutex );
    if( !flag )
      {
        printf("pong\n");
        flag = 1;
      }
    pthread_mutex_unlock( &mutex );
  } while(1);

  pthread_exit(0);
}

int main()
  {
     pthread_t ping, pong;
     pthread_mutex_init(&mutex, NULL);
     pthread_create( &ping, NULL, (void*)&ping_function, NULL );
     pthread_create( &pong, NULL, (void*)&pong_function, NULL );
     pthread_join( ping, NULL );
     pthread_join( pong, NULL );

     return 0;
  }



Zmienna warunku (ang. Conditional variable)
• mechanizm służący do synchronizacji działa wątków poprzez informowanie (ang.signaling) innych wątków o stanie współdzielonego zasobu (ang. shared resource)
• dopóki porządany stan wpółdzielonego zasobu nie zostanie osiągnięty wykonanie wątku może zostać zablokowane 

odblokowanie wątku może nastąpić przez inny wątek 
zawsze używana wraz z mutexem

pthread_cond_init(3p)                      inicjalizacja zmiennej warunku
pthread_cond_signal(3p)                  informowanie czekającego wątku, że określony warunek został spełniony
pthread_cond_broadcast(3p)            informowanie wszystkich czekających wątków, że określony warunek został osiągnięty
pthread_cond_wait(3p)                    czekanie, że określony warunek zostanie spełniony 

pthread_cond_destroy(3p)               zniszczenie zmiennej warunku





#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>  // EXIT_SUCCESS
#include <pthread.h> // pthread
#include <string.h>  // strcpy(), strcpy()

#define BUF_SIZE 100

char buff[BUF_SIZE];
int is_full = 0;
pthread_mutex_t my_mutex;
pthread_cond_t buff_cond_full, buff_cond_empty;

void* fun_write(void *arg)
{
        int i = 0;
        while(1) {
                pthread_mutex_lock(&my_mutex);

                        // jezeli nie ma miejsca do zapisu czekaj
                        if(is_full == 1)
                                pthread_cond_wait(&buff_cond_empty, &my_mutex);

                        if(i % 2) {
                                memcpy(buff, "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", BUF_SIZE);
                        } else {
                                memcpy(buff, "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", BUF_SIZE);
                        }
                        printf("WRITE index = %d\n", i++);

                        // oznacz, ze nie ma juz miejsca do zapisu i poinformuj o tym watek odczytujacy
                        is_full = 1;
                        pthread_cond_signal(&buff_cond_full);
                        
                pthread_mutex_unlock(&my_mutex);
        }
        pthread_exit(NULL);
}
void* fun_read(void* arg)
{
        char buff_local[BUF_SIZE];
        int index = 0;

        while(1) {
                pthread_mutex_lock(&my_mutex);

                        // jezeli nie ma danych czekaj
                        if(is_full == 0)
                                pthread_cond_wait(&buff_cond_full, &my_mutex);

                        memcpy(buff_local, buff, BUF_SIZE);                      
                        printf("READ index = %d | %s\n", index++, buff_local);

                        // oznacz, ze nie ma juz danych do odczytu i poinformuj watek zapisujacy
                        is_full = 0;
                        pthread_cond_signal(&buff_cond_empty);

                pthread_mutex_unlock(&my_mutex);
        }
        pthread_exit(NULL);
}

int main(int argc, char* argv[])
{
        pthread_t thread_write, thread_read;

        // inicjalizacja mutexu 
        pthread_mutex_init(&my_mutex, NULL);

        // inicjalizacja zmiennej warunku
        pthread_cond_init(&buff_cond_empty, NULL);
        pthread_cond_init(&buff_cond_full,  NULL);

        pthread_create(&thread_write, NULL, fun_write, NULL);
        pthread_create(&thread_read,  NULL, fun_read,  NULL);

        pthread_join(thread_write, NULL);
        pthread_join(thread_write, NULL);

        // zniszczenie mutexu
        pthread_mutex_destroy(&my_mutex);

        // zniszczenie zmiennej warunku 
        pthread_cond_destroy(&buff_cond_empty);
        pthread_cond_destroy(&buff_cond_full);

        exit(EXIT_SUCCESS);
}
przykłady M1/003thread_sync.../ 004


Semafor

•mechanizm służący do synchronizacji działania wątków/procesów

•posiadają wartość, która nie może spaść poniżej 0
•próba zmiejszenia wartości poniżej 0, powoduje zablokowanie działania
•inkrementacja/dekrementacja może nastąpić przez inny proces/wątek

nienazwane(ang.unnamed)
•nieposiadają nazwy-muszą znajdować siew pamięci wspólnej dla procesów/wątków
•mogą być używane dla:
      •wątków
      •procesów powiązanych

nazwane(ang.named)
•posiadają nazwe(/dev/shm) –nie muszą znajdować się we wspólnej pamięci
•mogą być używane dla:
     •wątki
     •procesy powiązane

     •procesy niepowiązane


Semafor -interfejs
sem_open(3) otwiera/tworzy nazwany semafor
sem_init(3) otwiera/tworzy nienazwany semafor

sem_post(3) inkrementacja semafora

sem_wait(3) dekrementacja semafora

sem_destroy(3) zniszczenie nienazwanego semafora

sem_close(3) zamknięcie nazwanego semafora


sem_unlink(3) usunięcie nazwanego semafora


Semafor nienazwany

#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>    // EXIT_SUCCESS
#include <pthread.h>   // pthread
#include <semaphore.h> // semaphore
#include <string.h>    // strcpy(), strcpy()

#define BUF_SIZE 100

char buff[BUF_SIZE];

// nienazwane semafory
sem_t sem_full;
sem_t sem_empty;

void* fun_write(void *arg)
{
 int i = 0;
 while(1) {
  sem_wait(&sem_empty); 
   // nie trzeba mutexu poniewaz jest tylko jeden dzielony zasob
   if(i % 2) {
    memcpy(buff, "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", BUF_SIZE);
   } else {
    memcpy(buff, "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", BUF_SIZE);
   }
   printf("WRITE index = %d\n", i++);
   
  sem_post(&sem_full);
 }
 pthread_exit(NULL);
}

void* fun_read(void* arg)
{
 char buff_local[BUF_SIZE];
 int index = 0;

 while(1) {
  sem_wait(&sem_full);
   // nie trzeba mutexu poniewaz jest tylko jeden dzielony zasob
   memcpy(buff_local, buff, BUF_SIZE);   
   printf("READ index = %d | %s\n", index++, buff_local);
   
  sem_post(&sem_empty);
   
 }
 pthread_exit(NULL);
}

int main(int argc, char* argv[])
{
 pthread_t thread_write, thread_read;
  
 // inicjalizacja semafora
 sem_init(
    &sem_full, // adres semafora 
       // sem_t *sem
        0,        // ma byc uzyty pomiedzy watkami
       // int pshared
     0   // poczatkowa wartosc
       // unsigned int value
    );

 // inicjalizacja semafora
 sem_init(&sem_empty, 0, 1); 
 
 pthread_create(&thread_write, NULL, fun_write, NULL);
 pthread_create(&thread_read,  NULL, fun_read,  NULL);

 pthread_join(thread_write, NULL);
 pthread_join(thread_read,  NULL);
 
 // zniszczenie semafora
 sem_destroy(&sem_full);
 sem_destroy(&sem_empty);
 
 exit(EXIT_SUCCESS);
}

przykłady M1/003thread_sync.../ 006

Semafor nazwany

#include <stdio.h>
#include <unistd.h>
#include <stdlib.h>    // EXIT_SUCCESS
#include <pthread.h>   // pthread
#include <semaphore.h> // semaphore
#include <fcntl.h>     // O_* constants
#include <string.h>    // strcpy(), strcpy()

#define BUF_SIZE 100

char buff[BUF_SIZE];

// nienazwane semafory
sem_t * sem_full  = NULL;
sem_t * sem_empty = NULL;

void* fun_write(void *arg)
{
 int i = 0;
 while(1) {
  sem_wait(sem_empty); 
   // nie trzeba mutexu poniewaz jest tylko jeden dzielony zasob
   if(i % 2) {
    memcpy(buff, "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", BUF_SIZE);
   } else {
    memcpy(buff, "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", BUF_SIZE);
   }
   printf("WRITE index = %d\n", i++);
   
  sem_post(sem_full);
 }
 pthread_exit(NULL);
}

void* fun_read(void* arg)
{
 char buff_local[BUF_SIZE];
 int index = 0;

 while(1) {
  sem_wait(sem_full);
   // nie trzeba mutexu poniewaz jest tylko jeden dzielony zasob
   memcpy(buff_local, buff, BUF_SIZE);   
   printf("READ index = %d | %s\n", index++, buff_local);
   
  sem_post(sem_empty);
   
 }
 pthread_exit(NULL);
}

int main(int argc, char* argv[])
{
 pthread_t thread_write, thread_read;

 // inicjalizacja semafora
 sem_full  = sem_open("/my_sem_full",  O_CREAT | O_EXCL | O_RDWR, 0600, 0); 

 // inicjalizacja semafora
 sem_empty = sem_open("/my_sem_empty", O_CREAT | O_EXCL | O_RDWR, 0600, 1); 
 
 pthread_create(&thread_write, NULL, fun_write, NULL);
 pthread_create(&thread_read,  NULL, fun_read, NULL);

 pthread_join(thread_write, NULL);
 pthread_join(thread_read,  NULL);
 
 // zamkniecie semaforow
 sem_close(sem_full);
 sem_close(sem_empty);
 
 // usuniecie semaforow
 sem_unlink("/my_sem_full");
 sem_unlink("/my_sem_empty");
 
 exit(EXIT_SUCCESS);
}

przykłady M1/003thread_sync.../ 007



KOMUNIKACJA MIĘDZYPROCESOWA

Potoki
     •potokinienazwane(ang.Pipes)
     •potokinazwane(ang.FIFOs)
 Kolejki komunikatów (ang.Message queues)
 Pamięć współdzielona (ang.Shared memory)

 Gniazda (ang.Sockets)



Potoki


potoki nienazwane(ang.pipes) oraz potoki nazwane(ang.FIFO) zapewniają jednokierunkową(ang.unidirectional) metodę komunikacji pomiędzy procesami polegającą na strumieniu bajtów(ang.byte stream)

•domyślnieodczyt(pusty) / zapis(pełny) jest blokujący

potoki nazwane
    •posiadają nazwę w systemie plików:
          •procesyniepowiązane/niespokrewnione(ang.unrelated)
          •Procesypowiązane/spokrewnione(ang.related)
potoki nienazwane można zleźć w /proc/<pid>/fd

•procesy powiązane/spokrewnione(ang.related)


Przykład

ls | sort –r
history | grepps

pipes unnamed_create
#include <stdio.h> 
#include <stdlib.h> // EXIT_SUCCESS
#include <string.h> // strlen
#include <unistd.h> // pipe, read, write

int main(int argc, char* argv[]) 
{ 
 int data_processed;
 int pipes_id[2]; 
  //pipes_id[0] - read
  //pipes_id[1] - write
 char buffer[32];
 
 if (pipe(pipes_id) == 0) {
  data_processed = write(pipes_id[1], "hello", strlen("hello"));
  printf("Wrote %d bytes\n", data_processed);
  
  data_processed = read(pipes_id[0], buffer, BUFSIZ);
  printf("Read %d bytes: %s\n", data_processed, buffer);
 }

 close(pipes_id[0]);
 close(pipes_id[1]);
 exit(EXIT_SUCCESS);
}

przykłady M2/002_IPC.../ 001


pipes unnamed_create.fork
#include <stdio.h>
#include <stdlib.h> // EXIT_SUCCESS
#include <string.h> // strlen
#include <unistd.h> // pipe, read, write

int main(int argc, char* argv[]) 
{ 
 int data_processed;
 int pid;
 int pipes_id[2]; 
  //pipes_id[0] - read
  //pipes_id[1] - write
 char buffer[32];
 
 if (pipe(pipes_id) == 0) {
  
  pid = fork();
  if(pid == -1) {  
  } else if (pid == 0) 
  { //child
   close(pipes_id[0]); // close read
   data_processed = write(pipes_id[1], "hello", strlen("hello"));
   printf("child wrote %d bytes\n", data_processed);
   close(pipes_id[1]);   
  } else 
  { //parent
   close(pipes_id[1]); // close write 
   data_processed = read(pipes_id[0], buffer, BUFSIZ);
   printf("parent read %d bytes: %s\n", data_processed, buffer);
   close(pipes_id[0]);
  }
 }
 
 exit(EXIT_SUCCESS);
}

przykłady M2/002_IPC.../ 001_pipes_unnamed/002


pipes_named_fifo_create

#include <stdio.h>
#include <stdlib.h>   // EXIT_SUCCESS
#include <sys/stat.h> // mkfifo

int main()
{ 
 // one descriptor : read / write
 int res = mkfifo("/tmp/example_fifo", 0777); 
 
 if (res == 0) 
  printf("FIFO created\n"); 

 exit(EXIT_SUCCESS); 
}

//one terminal
//cat < /tmp/example_fifo

//second terminal
//echo "hello" /tmp/example_fifo
przykłady M2/002_IPC.../ 002_pipes_named_fifo/001

pipes_named_fifo_comm
#include <stdio.h>
#include <string.h>    // strlen
#include <stdlib.h>    // EXIT_SUCCESS
#include <unistd.h>    // read,write
#include <sys/types.h> // open
#include <sys/stat.h>  // mkfifo
#include <fcntl.h>     // open flags

int main(int argc, char* argv[])
{ 
 int   fd;
 char  readbuf[255];
 
 if(argc < 2) {
  printf("incorrect usage\n");
  exit(EXIT_FAILURE);
 }
 
 if(0 == strcmp(argv[1], "1")) {
  printf("READ\n");
  
  mkfifo("/tmp/example_fifo", 0666);
  
  printf("before open\n");
  fd = open("/tmp/example_fifo", O_RDONLY);
  printf("after open\n");
  
  read(fd, readbuf, sizeof(readbuf));
  printf("read -  %s", readbuf);
  
  close(fd);
  unlink("/tmp/example_fifo"); // delete fifo
  exit(EXIT_SUCCESS);
  
 } else if(0 == strcmp(argv[1], "2")) {
  printf("WRITE\n");
  
     fd = open("/tmp/example_fifo", O_WRONLY);
  write(fd, "Hello world", strlen("Hello world"));
  close(fd);
  exit(EXIT_SUCCESS);
  
 } else {
  exit(EXIT_FAILURE);
 }
}



Kolejki komunikatów (message queues)

metoda komunikacji pomiędzy procesami oparta na wiadomościach(ang.message oriented communication)
•odbiorca odczytuje całe wiadomości, brak możliwości częściowego odczytu wiadomości
•wiadomości mają swoją strukturę i priorytet
      •0 -najniższy
      •31 –najwyższy(POSIX)
      •_SC_MQ_PRIO_MAX (32768) -najwyższy(Linux)



Pamięć współdzielona (Shared memory)

meteda komunikacji pomiędzy procesami oparta na współdzieleniu tego samego obszaru pamięci
•najszybsza metoda komunikacji ponieważ nie ma transferu danych:
•user space => kernel space => user space
•nie zapewnia synchronizacji
•tworzona w wirtualnym systemie plików: /dev/shm
•trzeba dolinkować bibliotekę real-time -lrt


Gniazda


najbardziej złożony mechanizm kumunikacji międzyprocesowej pozwalający również na komunikację zdalną (pomiędzy różnymi hostami)

•gniazdo(ang. socket) jest zakończeniem kanału kounikacyjnego
•potrzeba dwóch gniazd
•dwukierunkowe(ang. bidirectional)


Gniazda-interfejs

fd =  socket(domain, type, protocol)

domena(ang.domain):
     •metoda identyfikacji gniazda(format adresu)
     •zakres kumunikacji (lokalna, zdalna)
     •AF_UNIX/AF_LOCAL(UNIX_domain)
          •komunikacja w obrębie tego samego hosta
          •adres-ścieżka w systemie plików
     •AF_INET(IPv4 domain)
          •komunikacja poprzez sieć IPv4
          •adres–32bit IPv4 + numerportu
     •AF_INET6(IPv6) domain)
          •komunikacja poprzez siećIPv6
          •adres–128bit IPv6 + numerportu

typ(ang.type):
    •określasemantykękumunikacji(ang.semantics of communication)
    •SOCK_STREAM(ang.stream)
          •ciągbajtów(ang.byte stream) połączeniowa,niezawodna TCP
     SOCK_DGRAM(ang.datagram)
           •wiadomości(ang.message-oriented)
           •bezpołączeniowa(ang.connection-less), zawodna UDP


protokół(ang.protocol):
•określa konkretny protokół, kóry będzie użyty do komunikacji
•0 –domyślny protokół będzie używany(zdeterminowany przez domenę i typ)
       •AF_INET, SOCK_DGRAM => IPPROTO_UDP
       •AF_INET, SOCK_STREAM => IPPROTO_TCP
•!0 –można wybrać dowolny protokół
          •/etc/protocols => wszystkie dostępne protokoły



Gniazda-interfejs
socket(3) –tworzenie gniazda
bind(2) –łaczenie deskryptora z adresem lokalnym
connect(2) –łaczenie deskryptora gniazda z adresem zdalnym
listen(2) –ustawienie gniazda jako pasywnego(do akceptacji połączeń)
accept(3) –akceptacja połączenia na gnieździe
send(2), sendto(2), sendmsg(2) –wysyłaniedanych
rcv(2), rcvfrom(2), rcvmsg(2) –odbieranie danych
getsockname(2) –zwraca adres przywiązania do gniada
getpeername(2) –zwraca adres połaczonego klienta
getsockopt(2) –opcje gniazda
close(2), shutdown(2) –zamykanie gniazda





przykład local klient:


#include <stdio.h>
#include <stdlib.h>
#include <sys/un.h>     // sockaddr_un 
#include <sys/socket.h> // AF_UNIX
#include <unistd.h>     // sleep, close, unink

int main() 
{ 
 char   buff[20];
 int                client_sockfd;  
 struct sockaddr_un server_address; 
 
 // (1) stworz socket
 client_sockfd = socket(AF_UNIX,     // domain
         SOCK_STREAM, // type
         0);          // protocol
 
 // (2) polacz adres zdalny z desktryptorem gniazda, gniazdo aktywne
 server_address.sun_family = AF_UNIX; 
 strcpy(server_address.sun_path, "server_socket");
 connect(client_sockfd, (struct sockaddr *)&server_address, sizeof(server_address));

 //6) zapis
 sleep(5);
 write(client_sockfd, "hello world", strlen("hello world") + 1);
 
 //(5) odczyt, wywołanie blokujace
 read(client_sockfd, buff, 20);
 printf("Odczytano - %s\n", buff);
 
  // (7) - zamknięcie
 close(client_sockfd); 
}

local server


#include <stdio.h>
#include <stdlib.h>
#include <sys/un.h>     // sockaddr_un 
#include <sys/socket.h> // AF_UNIX
#include <unistd.h>     // sleep, close, unink

int main() 
{ 
 char   buff[20];
 int    server_sockfd;
 int    client_sockfd; 
 int    server_len; 
 int    client_len;
 struct sockaddr_un server_address;
 struct sockaddr_un client_address;
 
 // (1) stworz socket
 server_sockfd = socket(AF_UNIX,     // domain
         SOCK_STREAM, // type
         0);          // protocol
 
 // (2) połącz adres lokalny z desktryptorem gniazda
 server_address.sun_family = AF_UNIX;
 strcpy(server_address.sun_path, "server_socket");
 bind(server_sockfd, (struct sockaddr *)&server_address, sizeof(server_address)); 
 
 // (3) ustaw gniazdo jako pasywne
 listen(server_sockfd, 
        1); // The backlog argument defines the maximum length to which the queue of 
//pending connections for sockfd may grow   
 
 client_len = sizeof(client_address); 
 
 // (4) zaakceptuj polaczenie na gniezdzie, blokujace wywolanie
 printf("Czekamy aż klient się polaczy...\n");
 client_sockfd = accept(server_sockfd, (struct sockaddr *)&client_address, &client_len); 
 printf("Czekamy aż klient polaczyl sie\n");
 
 //(5) odczyt, wywolanie blokujace
 printf("Czekamy na dane od klienta...\n");
 read(client_sockfd, buff, 20);
 printf("Odczytano - %s\n", buff);
 
 //6) zapis
 write(client_sockfd, "hello back", strlen("hello back") + 1);
 
 sleep(5);
 
 // (7) - zamknięcie
 close(client_sockfd);
 unlink("server_socket");
}


network client


#include <stdio.h>
#include <stdlib.h>
#include <netinet/in.h> // sockaddr_in 
#include <arpa/inet.h>  // inet_addr
#include <sys/socket.h> // AF_UNIX
#include <unistd.h>     // sleep, close, unink
#include <string.h>     // strlen

int main() 
{ 
 char   buff[20];
 int                client_sockfd;  
 struct sockaddr_in server_address; 
 
 // (1) stworz socket
 client_sockfd = socket(AF_INET,     // domain
         SOCK_STREAM, // type
         0);          // protocol
 
 // (2) polacz adres zdalny z desktryptorem gniazda, gniazdo aktywne
 server_address.sin_family      = AF_INET; 
 server_address.sin_addr.s_addr = inet_addr("127.0.0.1");
 server_address.sin_port        = htons(9734); // => network order needed
 connect(client_sockfd, (struct sockaddr *)&server_address, sizeof(server_address));

 //6) zapis
 sleep(5);
 write(client_sockfd, "hello world", strlen("hello world") + 1);
 
 //(5) odczyt, wywołanie blokujace
 read(client_sockfd, buff, 20);
 printf("Odczytano - %s\n", buff);
 
  // (7) - zamknięcie
 close(client_sockfd); 
}


network server


#include <stdio.h>
#include <stdlib.h>
#include <netinet/in.h> // sockaddr_in 
#include <arpa/inet.h>  // inet_addr
#include <sys/socket.h> // AF_UNIX
#include <unistd.h>     // sleep, close, unink
#include <string.h>     // strlen

int main() 
{ 
 char   buff[20];
 int    server_sockfd;
 int    client_sockfd; 
 int    server_len; 
 int    client_len;
 struct sockaddr_in server_address;
 struct sockaddr_in client_address;
 
 // (1) stworz socket
 server_sockfd = socket(AF_INET,     // domain
         SOCK_STREAM, // type
         0);          // protocol
 
 // (2) połącz adres lokalny z desktryptorem gniazda
 server_address.sin_family      = AF_INET; 
 server_address.sin_addr.s_addr = inet_addr("127.0.0.1"); 
 server_address.sin_port        = htons(9734); // => network order needed
 bind(server_sockfd, (struct sockaddr *)&server_address, sizeof(server_address)); 
 
 // (3) ustaw gniazdo jako pasywne
 listen(server_sockfd, 
        1); // The backlog argument defines the maximum length to which the queue 
//of pending connections for sockfd may grow   
 
 client_len = sizeof(client_address); 
 
 // (4) zaakceptuj polaczenie na gniezdzie, blokujace wywolanie
 printf("Czekamy aż klient się polaczy...\n");
 client_sockfd = accept(server_sockfd, (struct sockaddr *)&client_address, &client_len); 
 printf("Czekamy aż klient polaczyl sie\n");
 
 //(5) odczyt, wywolanie blokujace
 printf("Czekamy na dane od klienta...\n");
 read(client_sockfd, buff, 20);
 printf("Odczytano - %s\n", buff);
 
 //6) zapis
 write(client_sockfd, "hello back", strlen("hello back") + 1);
 
 sleep(5);
 
 // (7) - zamknięcie
 close(client_sockfd);
}
.


Atrybuty Schedulera

Scheduling Policy  SCHED_FIFO  i SCHED_RR

przykład:

#include<stdlib.h>
#include<sched.h>
#include<stdio.h>
#include<pthread.h>
#include<unistd.h>  // sleep

                                    //Scheduling policy
pthread_attr_t my_attr;                                          //atrybuty threadow
pthread_attr_t my_attr2;                                          //atrybuty threadoww
struct sched_param param1, param2;

long unsigned int ff1 =0;
long unsigned int ff2=0;

void funkcja1(void)     
{
while(1)                      
{
    if (!(ff1 % 10000000)) {
    printf("funkcja1 = %lu\n",ff1);
    }
    ff1++;
}
}

void funkcja2(void)     
{
while(1)                
{
    if (!(ff2 % 10000000)) {
    printf("funkcja2 = %lu\n",ff2);
    }
    ff2++;
}
}

int main ()
{
pthread_attr_init(&my_attr);
pthread_attr_init(&my_attr2);
pthread_attr_setinheritsched (&my_attr, PTHREAD_EXPLICIT_SCHED);     //nie dziedzicz ustawien watku z wywalujacego
pthread_attr_setinheritsched (&my_attr2, PTHREAD_EXPLICIT_SCHED);  
pthread_attr_setschedpolicy(&my_attr, SCHED_RR);  //RR lub FIFO
pthread_attr_setschedpolicy(&my_attr2, SCHED_RR); 
param1.sched_priority = 1;
param2.sched_priority = 5;

pthread_attr_setschedparam(&my_attr, &param1);               //param1 do atrybutow  watku
pthread_t ping;                     //tworzenie watku 1
pthread_create( &ping, &my_attr , (void*)&funkcja1, NULL );   //uruchamianie funkcji1  

pthread_attr_setschedparam(&my_attr2, &param2);              //param2 do atrybutow watku
pthread_t pong;                     //tworzenie watku 2
pthread_create( &pong, &my_attr,  (void*)&funkcja2, NULL );   //uruchamianie funkcji2

pthread_join( ping, NULL );    //czeka az watek pink zakonczy
pthread_join( pong, NULL );    //czeka az watek pink zakonczy
pthread_attr_destroy(&my_attr);   //niszczy zmienna atrybutow
pthread_attr_destroy(&my_attr2);   //niszczy zmienna atrybutow
return 0;
}

nie są konieczne dwie zmienne atrybutów (my_attr i my_attr2 )

Linki:
pierwszy
drugi
trzeci
manual

..