Method and apparatus for reaching agreement between nodes in a distributed system
22 claims: 22 independent, 0 dependent
- 1A method for selecting a node to host a primary server (106) for a service (202,203,204,205) from a plurality of candidate nodes (102,103,104,105) in a distributed computing system (100), the method comprising:a) receiving (401) an indication that a state of the distributed computing system has changed;b) in response to the indication, determining (602) if a node that was previously hosting the primary server for the service (106) continues to exist;c) if there is not already a node hosting the primary server for the service, selecting (604) a new node to host the primary server based upon rank information for the candidate nodes by communicating rank information between a given node and other nodes in the distributed computing system, wherein each node in the distributed computing system has a unique rank with respect to the other nodes in the distributed computing system, comparing a rank of the given node with a rank of the other nodes in the distributed computing system, and if one of the other nodes in the distributed computing system has a higher rank than the given node, disqualifying (612) the given node from hosting the primary server;characterized by: d) periodically sending (502) checkpointing information (120,121) from the primary server (106) to at least one secondary server (107,108) to maintain a consistent state with the primary server for the same service, whereby a secondary server is able to take over from the primary server if the primary server fails or otherwise becomes unavailable;e) enabling a given node (I) from the plurality of nodes (102,103,104,105) in the distributed computing system to act as one of: - a host node (102) for the primary server (106) for the service;or- a host node (103,104) for a secondary server (107, 108) for the service which has received the checkpointing information (120,121) from the primary server;or- a spare node (105) for the primary server, wherein the spare node does not receive checkpointing information from the primary server,f) wherein selecting (604) a new node to host the primary server for the service is performed concurrently on all active nodes by means of distributed selection mechanisms (132-135) communicating rank information with each other through a set of shared, lockable candidate variables which contain an identifier for a candidate node to host the primary server, andg) wherein disqualifying (612) the given node (I) from hosting the primary server comprises writing a new identifier into the candidate variable for the given node (I) if the rank of the given node is less than the rank of the new node. Procédé pour sélectionner un noeud pour l'hébergement d'un serveur primaire (106) pour un service (202, 203, 204, 205) parmi une pluralité de noeuds candidats (102, 103, 104, 105) dans un système informatique distribué (100), le procédé consistant à : a) recevoir (401) une indication selon laquelle un état du système informatique distribué a changé ;b) en réponse à l'indication, déterminer (602) si un noeud qui a précédemment hébergé le serveur primaire pour le service (106) existe toujours ;c) s'il n'y a pas déjà de noeud hébergeant le serveur primaire pour le service, sélectionner (604) un nouveau noeud pour l'hébergement du serveur primaire sur la base d'informations de rang pour les noeuds candidats en communiquant des informations de rang entre un noeud donné et d'autres noeuds dans le système informatique distribué, chaque noeud dans le système informatique distribué ayant un rang unique relativement aux autres noeuds dans le système informatique distribué, et comparer un rang du noeud donné avec un rang des autres noeuds dans le système informatique distribué, et si l'un des autres noeuds dans le système informatique distribué a un rang supérieur à celui du noeud donné, disqualifier (612) le noeud donné en tant que noeud d'hébergement du serveur primaire;caractérisé par : d) l'envoi périodique (502) d'informations de point de contrôle (120, 121) du serveur primaire (106) à au moins un serveur secondaire (107, 108) pour maintenir un état de concordance avec le serveur primaire pour le même service, si bien qu'un serveur secondaire est capable de prendre le relais du serveur primaire si le serveur primaire subit une défaillance ou d'une autre manière devient indisponible ;e) la permission à un noeud donné (I) parmi la pluralité de noeuds (102, 103, 104, 105) dans le système informatique distribué de servir de l'un de - un noeud hôte (102) pour le serveur primaire (106) pour le service ;ou- un noeud hôte (103, 104) pour un serveur secondaire (107, 108) pour le service qui a reçu les informations de point de contrôle (120, 121) du serveur primaire;ou- un noeud de réserve (105) pour le serveur primaire, le noeud de réserve ne recevant pas d'informations de point de contrôle du serveur primaire,f) la sélection (604) d'un nouveau noeud pour l'hébergement du serveur primaire pour le service étant exécutée concurremment sur tous les noeuds actifs à l'aide de mécanismes de sélection distribués (132 à 135) communiquant des informations de rang entre eux par l'intermédiaire d'un ensemble de variables candidat verrouillables, partagées, qui contiennent un identificateur pour un noeud candidat pour l'hébergement du serveur primaire, etg) la disqualification (612) du noeud donné en tant que noeud d'hébergement du serveur primaire consistant à écrire un nouvel identificateur dans la variable candidat pour le noeud donné (I) si le rang du noeud donné est inférieur au rang du nouveau noeud. Verfahren, um aus mehreren Knotenkandidaten (102, 103, 104, 105) in einem verteilten Computersystem (100) einen Knoten auszuwählen, damit er das Hosting eines primären Servers (106) für einen Dienst (202, 203, 204, 205) ausführt, wobei das Verfahren umfasst: a) Empfangen (401) eines Hinweises, dass sich ein Zustand des verteilten Computersystems geändert hat;b) in Reaktion auf den Hinweis Bestimmen (602), ob ein Knoten, der früher das Hosting des primären Servers für den Dienst (106) ausgeführt hat, weiterhin existiert;c) falls nicht bereits ein Knoten vorhanden ist, der das Hosting des primären Servers für den Dienst ausführt, Auswählen (604) eines neuen Knotens, um das Hosting des primären Servers auszuführen, auf der Grundlage von Ranginformationen für die Knotenkandidaten, indem die Ranginformationen zwischen einem gegebenen Knoten und anderen Knoten in dem verteilten Computersystem ausgetauscht werden, wobei jeder Knoten in dem verteilten Computersystem einen eindeutigen Rang in Bezug auf den anderen Knoten in dem verteilten Computersystem hat, ein Rang des gegebenen Knotens mit einem Rang der anderen Knoten in dem verteilten Computersystem verglichen wird und dann, falls einer der anderen Knoten in dem verteilten Computersystem einen höheren Rang als der gegebene Knoten hat, der gegebene Knoten von der Ausführung des Hostings des primären Servers ausgeschlossen wird (612);gekennzeichnet durch: d) periodisches Senden (502) von Fixpunktroutinen-Informationen (120, 121) von dem primären Server (106) zu wenigstens einem sekundären Server (107, 108), um einen konsistenten Zustand mit dem primären Server für denselben Dienst aufrecht zu erhalten, wobei der sekundäre Server von dem primären Server übernehmen kann, falls der primäre Server ausfällt oder auf andere Weise unverfügbar wird;e) Freigeben eines gegebenen Knotens (I) unter den mehreren Knoten (102, 103, 104, 105) in dem verteilten Computersystem, damit er wirkt als: - ein Host-Knoten (102) für den primären Server (106) für den Dienst;oder- ein Host-Knoten (103, 104) für einen sekundären Server (107, 108) für den Dienst, der die Fixpunktroutinen-Informationen (120, 121) für den primären Server empfangen hat;oder- ein Reserveknoten (105) für den primären Server, wobei der Reserveknoten Fixpunktroutinen-Informationen von dem primären Server nicht empfängt,f) wobei das Auswählen (604) eines neuen Knotens, um das Hosting des primären Servers für den Dienst auszuführen, konkurrent auf allen aktiven Knoten mittels verteilter Auswahlmechanismen (132-135), die miteinander Ranginformationen über eine Menge gemeinsam genutzter, fixierbarer Variablenkandidaten, die einen Identifizierer für einen Knotenkandidaten für die Ausführung des Hostings des primären Servers enthalten, austauschen, ausgeführt wird, undg) wobei das Ausschließen (612) des gegebenen Knotens (I) von der Ausführung des Hostings des primären Servers das Schreiben eines neuen Identifizierers in den Variablenkandidaten für den gegebenen Knoten (I) umfasst, falls der Rang des gegebenen Knotens niedriger als der Rang des neuen Knotens ist.
- 2Procédé selon la revendication 1, consistant en outre, s'il existe toujours un noeud qui est configuré pour héberger le serveur primaire, à permettre (622) au noeud qui est configuré pour héberger le serveur primaire de communiquer avec d'autres noeuds dans le système informatique distribué de manière à disqualifier (624) les autres noeuds en tant que noeuds d'hébergement du serveur primaire. The method of claim 1, further comprising, if there continues to exist a node that is configured to host the primary server, allowing (622) the node that is configured to host the primary server to communicate with other nodes in the distributed computing system in order to disqualify (624) the other nodes from hosting the primary server. Verfahren nach Anspruch 1, das ferner dann, wenn ein Knoten, der so konfiguriert ist, dass er das Hosting des primären Servers ausführt, weiterhin existiert, das Zulassen (622), dass der Knoten, der so konfiguriert ist, dass er das Hosting des primären Servers ausführt, mit anderen Knoten in dem verteilten Computersystem kommuniziert, umfasst, um die anderen Knoten von der Ausführung des Hostings des primären Servers auszuschließen (624).
- 3Procédé selon la revendication 1, consistant en outre à fixer initialement (700) la variable candidat pour identifier le noeud donné (I). The method of claim 1, further comprising, initially setting (700) the candidate variable to identify the given node (I). Verfahren nach Anspruch 1, das ferner das anfängliche Setzen (700) des Variablenkandidaten umfasst, um den gegebenen Knoten (I) zu identifizieren.
- 4Procédé selon la revendication 1, consistant en outre, après qu'un nouveau noeud a été sélectionné pour héberger le serveur primaire (106), si le nouveau noeud est différent d'un noeud précédent qui a hébergé le serveur primaire, à établir (108) des connexions au nouveau noeud pour le service. The method of claim 1, further comprising, after a new node has been selected to host the primary server (106), if the new node is different from a previous node that hosted the primary server, establishing (408) connections for the service to the new node. Verfahren nach Anspruch 1, das ferner nach der Auswahl eines neuen Knotens für die Ausführung des Hostings des primären Servers (106) dann, wenn der neue Knoten von einem früheren Knoten verschieden ist, der das Hosting des primären Servers ausgeführt hat, das Herstellen (408) von Verbindungen für den Dienst zu dem neuen Knoten umfasst.
- 5Procédé selon la revendication 1, consistant en outre, après qu'un nouveau noeud a été sélectionné pour héberger le serveur primaire (106), si le nouveau noeud est différent d'un noeud précédent qui a hébergé le serveur primaire, à configurer (410) le nouveau noeud pour l'hébergement du serveur primaire pour le service. The method of claim 1, further comprising, after a new node has been selected to host the primary server (106), if the new node is different from a previous node that hosted the primary server, configuring (410) the new node to host the primary server for the service. Verfahren nach Anspruch 1, das ferner nach der Auswahl eines neuen Knotens für die Ausführung des Hostings des primären Servers (106) dann, wenn der neue Knoten von einem früheren Knoten verschieden ist, der das Hosting des primären Servers ausgeführt hat, das Konfigurieren (410) des neuen Knotens für die Ausführung des Hostings des primären Servers für den Dienst umfasst.
- 6Procédé selon la revendication 1, consistant en outre à redémarrer (412) le service si le service a été interrompu suite au changement de l'état du système informatique distribué. The method of claim 1, further comprising restarting (412) the service if the service was interrupted as a result of the change in state of the distributed computing system. Verfahren nach Anspruch 1, das ferner das Neustarten (412) des Dienstes umfasst, falls der Dienst als Ergebnis der Änderung des Zustands des verteilten Computersystems unterbrochen wurde.
- 7Procédé selon la revendication 1, consistant en outre, lors du démarrage initial du service, à sélectionner un noeud de réserve de rang le plus élevé pour héberger le serveur primaire (106) pour le service. The method of claim 1, further comprising, upon initial startup of the service, selecting a highest ranking spare node to host the primary server (106) for the service. Verfahren nach Anspruch 1, das ferner bei einem anfänglichen Einrichten des Dienstes das Auswählen eines Reserveknotens mit höchstem Rang für die Ausführung des Hostings des primären Servers (106) für den Dienst umfasst.
- 8Procédé selon la revendication 1, consistant en outre à permettre au serveur primaire (106) de favoriser (504) des noeuds de réserve (105) dans le système informatique distribué pour l'hébergement des serveurs secondaires (107, 108) pour le service. The method of claim 1, further comprising allowing the primary server (106) to promote (504) spare nodes (105) in the distributed computing system to host secondary servers (107,108) for the service. Verfahren nach Anspruch 1, das ferner das Zulassen, dass der primäre Server (106) Reserveknoten (105) in dem verteilten Computersystem bei der Ausführung des Hostings sekundärer Server (107, 108) für den Dienst begünstigt, umfasst.
- 9Procédé selon la revendication 1, dans lequel la comparaison du rang du noeud donné (I) avec le rang des autres noeuds dans le système informatique distribué implique de considérer qu'un hôte pour le serveur primaire (106) a un rang supérieur à celui d'un hôte pour un serveur secondaire (107, 108), et de considérer qu'un hôte pour un serveur secondaire (107, 108) a un rang supérieur à celui d'un noeud de réserve (105). The method of claim 1, wherein comparing the rank of the given node (I) with the rank of the other nodes in the distributed computing system involves considering a host for the primary server (106) to have a higher rank than a host for a secondary server (107,108), and considering a host for a secondary server (107,108) to have a higher rank than a spare (105). Verfahren nach Anspruch 1, bei dem das Vergleichen des Rangs des gegebenen Knotens (I) mit dem Rang der anderen Knoten in dem verteilten Computersystem das Betrachten eines Hosts für den primären Server (106), der einen höheren Rang als ein Host für einen sekundären Server (107, 108) besitzt, und das Betrachten eines Hosts für einen sekundären Server (107, 108), der einen höheren Rang als ein Reserveknoten (105) besitzt, umfasst.
- 10Procédé selon la revendication 1, dans lequel la disqualification (612) du noeud donné (I) en tant que noeud d'hébergement du serveur primaire (106) implique de cesser (614) de communiquer des informations de rang entre le noeud donné et les autres noeuds dans le système informatique distribué. The method of claim 1, wherein disqualifying (612) the given node (I) from hosting the primary server (106) involves ceasing (614) to communicate rank information between the given node and the other nodes in the distributed computing system. Verfahren nach Anspruch 1, bei dem das Ausschließen (612) des gegebenen Knotens (I) von der Ausführung des Hostings des primären Servers (106) das Beenden (614) des Austauschens von Ranginformationen zwischen dem gegebenen Knoten und den anderen Knoten in dem verteilten Computersystem umfasst.
- 11A computer program (134), which when executing on a distributed computer network (110), performs the method steps of any one of claims 1 to 10. Computerprogramm (134), das dann, wenn es in einem verteilten Computernetz (110) ausgeführt wird, die Verfahrensschritte nach einem der Ansprüche 1 bis 10 ausführt. Programme informatique (134), qui, lorsqu'il est exécuté sur un réseau informatique distribué (110), exécute les étapes de procédé de l'une quelconque des revendications 1 à 10.
- 12A computer-readable storage medium storing the computer program (134) of claim 11. Computerlesbares Speichermedium, das das Computerprogramm (134) nach Anspruch 11 speichert. Support de stockage lisible par ordinateur stockant le programme informatique (134) de la revendication 11.
- 13An apparatus (100) for selecting a node to host a primary server (106) for a service (202,203,204,205) from a plurality of candidate nodes (102,103,104,105) in a distributed computing system (100), the apparatus comprising:a) a receiving mechanism that is configured to receive (401) an indication that a state of the distributed computing system has changed;b) a determining mechanism that is configured to determine (602) if a node that was previously hosting the primary server for the service (106) continues to exist, in response to the indication;c) a distributed selection selection mechanism (132-135) at each one of the plurality of nodes (102,103,104,105) in communication with each other that is configured to select (604) a new node to host the primary server based upon rank information for the candidate nodes, if there is not already a node hosting the primary server for the service;by means of a communicating mechanism configured to communicate rank information between a given node and other nodes in the distributed computing system, wherein each node in the distributed computing system has a unique rank with respect to the other nodes in the distributed computing system;a comparing mechanism configured to compare a rank of the given node with a rank of the other nodes in the distributed computing system;anda disqualifying mechanism configured to disqualify (612) the given node from hosting the primary server, if one of the other nodes in the distributed computing system has a higher rank than the given node;characterized by: d) means for periodically sending (502) checkpointing information (120,121) from the primary server (106) to at least one secondary server (107,108) to maintain a consistent state with the primary server for the same service, whereby a secondary server is able to take over from the primary server if the primary server fails or otherwise becomes unavailable;e) means for enabling a given node (I) from the plurality of nodes (102,103,104,105) in the distributed computing system to act as one of: - a host node (102) for the primary server (106) for the service;or- a host node (103,104) for a secondary server (107, 108) for the service which has received the checkpointing information (120,121) from the primary server;or- a spare node (105) for the primary server, wherein the spare node does not receive checkpointing information from the primary server,f) wherein the distributed selection mechanism (132-135) is configured to select (604) a new node to host the primary server for the service concurrently on all active nodes by means of communicating rank information with each other through a set of shared, lockable candidate variables which contain an identifier for a candidate node to host the primary server, andg) wherein the disqualifying mechanism comprises means for disqualifying (612) the given node (I) from hosting the primary server by writing a new identifier into the candidate variable for the given node (I) if the rank of the given node is less than the rank of the new node. Appareil (100) pour sélectionner un noeud pour l'hébergement d'un serveur primaire (106) pour un service (202, 203, 204, 205) parmi une pluralité de noeuds candidats (102, 103, 104, 105) dans un système informatique distribué (100), l'appareil comprenant : a) un mécanisme de réception qui est configuré pour recevoir (401) une indication selon laquelle un état du système informatique distribué a changé ;b) un mécanisme de détermination qui est configuré pour déterminer (602) si un noeud qui a précédemment hébergé le serveur primaire pour le service (106) existe toujours, en réponse à l'indication ;c) un mécanisme de sélection distribué (132 à 135) au niveau de chaque noeud de la pluralité de noeuds (102, 103, 104, 105) communiquant entre eux qui est configuré pour sélectionner (604) un nouveau noeud pour l'hébergement du serveur primaire sur la base d'informations de rang pour les noeuds candidats, s'il n'y a pas déjà un noeud hébergeant le serveur primaire pour le service ;à l'aide d'un mécanisme de communication configuré pour communiquer des informations de rang entre un noeud donné et d'autres noeuds dans le système informatique distribué, chaque noeud dans le système informatique distribué ayant un rang unique relativement aux autres noeuds dans le système informatique distribué ;d'un mécanisme de comparaison configuré pour comparer un rang du noeud donné avec un rang des autres noeuds dans le système informatique distribué ;et d'un mécanisme de disqualification configuré pour disqualifier (612) le noeud donné en tant que noeud d'hébergement du serveur primaire, si l'un des autres noeuds dans le système informatique distribué a un rang supérieur à celui du noeud donné ;caractérisé par : d) des moyens pour envoyer périodiquement (502) des informations de point de contrôle (120, 121) du serveur primaire (106) à au moins un serveur secondaire (107, 108) de manière à maintenir un état de concordance avec le serveur primaire pour le même service, si bien qu'un serveur secondaire est capable de prendre le relais du serveur primaire si le serveur primaire subit une défaillance ou d'une autre manière devient indisponible ;e) des moyens pour permettre à un noeud donné (I) parmi la pluralité de noeuds (102, 103, 104, 105) dans le système informatique distribué de servir de l'un de : - un noeud hôte (102) pour le serveur primaire (106) pour le service ;ou- un noeud hôte (103, 104) pour un serveur secondaire (107, 108) pour le service qui a reçu du serveur primaire les informations de point de contrôle (120, 121) ;ou- un noeud de réserve (105) pour le serveur primaire, le noeud de réserve ne recevant pas d'informations de point de contrôle du serveur primaire,f) le mécanisme de sélection distribué (132 à 135) étant configuré pour sélectionner (604) un nouveau noeud pour l'hébergement du serveur primaire pour le service concurremment sur tous les noeuds actifs au moyen de la communication d'informations de rang entre eux par l'intermédiaire d'un ensemble de variables candidat verrouillables, partagées, qui contiennent un identificateur pour un noeud candidat pour l'hébergement du serveur primaire, etg) le mécanisme de disqualification comprenant des moyens pour disqualifier (612) le noeud donné (I) en tant que noeud d'hébergement du serveur primaire en écrivant un nouvel identificateur dans la variable candidat pour le noeud donné (I) si le rang du noeud donné est inférieur au rang du nouveau noeud. Vorrichtung (100), um aus mehreren Knotenkandidaten (102, 103, 104, 105) in einem verteilten Computersystem (100) einen Knoten auszuwählen, damit er das Hosting eines primären Servers (106) für einen Dienst (202, 203, 204, 205) ausführt, wobei die Vorrichtung umfasst: a) einen Empfangsmechanismus, der so konfiguriert ist, dass er einen Hinweis darüber empfängt (401), dass sich ein Zustand des verteilten Computersystems geändert hat;b) einen Bestimmungsmechanismus, der so konfiguriert ist, dass er in Reaktion auf den Hinweis bestimmt (602), ob ein Knoten, der früher das Hosting des primären Servers für den Dienst (106) ausgeführt hat, weiterhin existiert;c) einen verteilten Auswahlmechanismus (132-135) bei jedem der mehreren Knoten (102, 103, 104, 105), die miteinander kommunizieren, der so konfiguriert ist, dass er einen neuen Knoten für die Ausführung des Hostings des primären Servers anhand von Ranginformationen für die Knotenkandidaten auswählt (604), falls nicht bereits ein Knoten das Hosting des primären Servers für den Dienst ausführt;mittels eines Kommunikationsmechanismus, der so konfiguriert ist, dass er Ranginformationen zwischen einem gegebenen Knoten und anderen Knoten in dem verteilten Computersystem austauscht, wobei jeder Knoten in dem verteilten Computersystem einen eindeutigen Rang in Bezug auf die anderen Knoten in dem verteilten Computersystem hat;eines Vergleichsmechanismus, der so konfiguriert ist, dass er einen Rang des gegebenen Knotens mit einem Rang der anderen Knoten in dem verteilten Computersystem vergleicht;und eines Ausschließungsmechanismus, der so konfiguriert ist, dass er den gegebenen Knoten von der Ausführung des Hostings des primären Servers ausschließt (612), falls einer der anderen Knoten in dem verteilten Computersystem einen höheren Rang als der gegebene Knoten hat;gekennzeichnet durch: d) eine Einrichtung zum periodischen Senden (502) von Fixpunktroutinen-Informationen (120, 121) von dem primären Server (106) zu wenigstens einem sekundären Server (107, 108), um einen konsistenten Zustand mit dem primären Server für denselben Dienst aufrecht zu erhalten, wobei ein sekundärer Server von dem primären Server übernehmen kann, falls der primäre Server ausfällt oder auf andere Weise unverfügbar wird;e) eine Einrichtung zum Freigeben eines gegebenen Knotens (I) unter den mehreren Knoten (102, 103, 104, 105) in dem verteilten Computersystem, damit er wirkt als: - ein Host-Knoten (102) für den primären Server (106) für den Dienst;oder- ein Host-Knoten (103, 104) für einen sekundären Server (107, 108) für den Dienst, der die Fixpunktroutinen-Informationen (120, 121) von dem primären Server empfangen hat;oder- ein Reserveknoten (105) für den primären Server, wobei der Reserveknoten keine Fixpunktroutinen-Informationen von dem primären Server empfängt,f) wobei der verteilte Auswahlmechanismus (132-135) so konfiguriert ist, dass er einen neuen Knoten für die Ausführung des Hostings des primären Servers für den Dienst konkurrent bei allen aktiven Knoten mittels des Austauschens von Ranginformationen untereinander über eine Menge gemeinsam genutzter, fixierbarer Variablenkandidaten, die einen Identifizierer für einen Knotenkandidaten für die Ausführung des Hostings des primären Servers enthalten, auswählt (604), undg) wobei der Ausschließungsmechanismus eine Einrichtung umfasst, um den gegebenen Knoten (I) von der Ausführung des Hostings des primären Servers auszuschließen, indem sie einen neuen Identifizierer in den Variablenkandidaten für den gegebenen Knoten (I) schreibt, falls der Rang des gegebenen Knotens niedriger als der Rang des neuen Knotens ist.
- 14Appareil selon la revendication 13, comprenant en outre un mécanisme qui est configuré pour permettre (622) au noeud qui est configuré pour héberger le serveur primaire de communiquer avec d'autres noeuds dans le système informatique distribué de manière à disqualifier (624) les autres noeuds en tant que noeuds d'hébergement du serveur primaire, s'il existe toujours un noeud qui est configuré pour héberger le serveur primaire. The apparatus of claim 13, further comprising a mechanism that is configured to allow (622) the node that is configured to host the primary server to communicate with other nodes in the distributed computing system in order to disqualify (624) the other nodes from hosting the primary server, if there continues to exist a node that is configured to host the primary server. Vorrichtung nach Anspruch 13, die ferner einen Mechanismus umfasst, der so konfiguriert ist, dass er dem Knoten, der so konfiguriert ist, dass er das Hosting des primären Servers ausführt, erlaubt (622) mit anderen Knoten in dem verteilten Computersystem zu kommunizieren, um die anderen Knoten von der Ausführung des Hostings des primären Servers auszuschließen, falls ein Knoten, der so konfiguriert ist, dass er das Hosting des primären Servers ausführt, weiterhin existiert.
- 15Appareil selon la revendication 13, dans lequel le mécanisme de sélection est configuré pour fixer initialement la variable candidat pour identifier le noeud donné (I). The apparatus of claim 13, wherein the selecting mechanism is configured to initially set the candidate variable to identify the given node (I). Vorrichtung nach Anspruch 13, bei der der Auswahlmechanismus so konfiguriert ist, dass er anfangs den Variablenkandidaten setzt, um den gegebenen Knoten (I) zu identifizieren.
- 16Appareil selon la revendication 13, comprenant en outre un mécanisme de connexion qui est configuré pour établir (408) des connexions à un nouveau noeud pour le service, après que le nouveau noeud a été sélectionné pour héberger le serveur primaire (106), et si le nouveau noeud est différent d'un noeud précédent qui a hébergé le serveur primaire. The apparatus of claim 13, further comprising a connection mechanism that is configured to establish (408) connections for the service to a new node, after the new node has been selected to host the primary server (106), and if the new node is different from a previous node that hosted the primary server. Vorrichtung nach Anspruch 13, die ferner einen Verbindungsmechanismus umfasst, der so konfiguriert ist, dass er Verbindungen für den Dienst zu einem neuen Knoten herstellt (408), nachdem der neue Knoten ausgewählt worden ist, das Hosting des primären Servers (106) auszuführen, und falls der neue Knoten von einem früheren Knoten, der das Hosting des primären Servers ausgeführt hat, verschieden ist.
- 17Appareil selon la revendication 13, comprenant en outre un mécanisme qui est configuré pour configurer (410) un nouveau noeud pour l'hébergement du serveur primaire pour le service, après que le nouveau noeud a été sélectionné pour héberger le serveur primaire (106), et si le nouveau noeud est différent d'un noeud précédent qui a hébergé le serveur primaire. The apparatus of claim 13, further comprising a mechanism that is configured to configure (410) a new node to host the primary server for the service, after the new node has been selected to host the primary server (106), and if the new node is different from a previous node that hosted the primary server. Vorrichtung nach Anspruch 13, die ferner einen Mechanismus umfasst, der so konfiguriert ist, dass er einen neuen Knoten konfiguriert (410), um das Hosting des primären Servers für den Dienst auszuführen, nachdem der neue Knoten für die Ausführung des Hostings des primären Servers (106) ausgewählt worden ist und falls der neue Knoten von einem früheren Knoten, der das Hosting des primären Servers ausgeführt hat, verschieden ist.
- 18Appareil selon la revendication 13, comprenant en outre un mécanisme de redémarrage (412) qui est configuré pour redémarrer le service si le service a été interrompu suite au changement de l'état du système informatique distribué (100). The apparatus of claim 13, further comprising a restarting mechanism (412) that is configured to restart the service if the service was interrupted as a result of the change in state of the distributed computing system (100). Vorrichtung nach Anspruch 13, die ferner einen Neustartmechanismus (412) umfasst, der so konfiguriert ist, dass er den Dienst neu startet, falls der Dienst als Ergebnis der Zustandsänderung des verteilten Computersystems (100) unterbrochen wurde.
- 19Appareil selon la revendication 13, comprenant en outre un mécanisme d'initialisation, dans lequel, pendant l'initialisation du service, le mécanisme d'initialisation est configuré pour sélectionner un noeud de réserve de rang le plus élevé pour l'hébergement du serveur primaire (106) pour le service. The apparatus of claim 13, further comprising an initialization mechanism wherein during initialization of the service, the initialization mechanism is configured to select a highest-ranking spare node to host the primary server (106) for the service. Vorrichtung nach Anspruch 13, die ferner einen Initialisierungsmechanismus umfasst, wobei während der Initialisierung des Dienstes der Initialisierungsmechanismus so konfiguriert ist, dass er einen Reserveknoten mit höchstem Rang auswählt, um das Hosting des primären Servers (106) für den Dienst auszuwählen.
- 20Appareil selon la revendication 13, comprenant en outre un mécanisme de favorisation qui est configuré pour favoriser (504) des noeuds de réserve dans le système informatique distribué (100) pour l'hébergement des serveurs secondaires (107, 108) pour le service. The apparatus of claim 13, further comprising a promotion mechanism that that is configured to promote (504) spare nodes in the distributed computing system (100) to host secondary servers (107,108) for the service. Vorrichtung nach Anspruch 13, die ferner einen Begünstigungsmechanismus umfasst, der so konfiguriert ist, dass er Reserveknoten in dem verteilten Computersystem (100) bei der Ausführung des Hostings sekundärer Server (107, 108) für den Dienst begünstigt.
- 21Appareil selon la revendication 13, dans lequel, tandis qu'est effectuée la comparaison du rang du noeud donné (I) avec le rang des autres noeuds dans le système informatique distribué, le mécanisme de sélection est configuré pour considérer qu'un hôte pour le serveur primaire (106) a un rang supérieur à celui d'un hôte pour un serveur secondaire (107, 108), et pour considérer qu'un hôte pour un serveur primaire a un rang supérieur à celui d'un noeud de réserve (105). The apparatus of claim 13, wherein while comparing the rank of the given node (I) with the rank of the other nodes in the distributed computing system, the selecting mechanism is configured to consider a host for the primary server (106) to have a higher rank than a host for a secondary server (107,108), and to consider a host for a secondary server to have a higher rank than a spare (105). Vorrichtung nach Anspruch 13, bei der während des Vergleichens des Rangs des gegebenen Knotens (I) mit dem Rang der anderen Knoten in dem verteilten Computersystem der Auswahlmechanismus so konfiguriert ist, dass er einen Host für den primären Server (106) betrachtet, der einen höheren Rang als ein Host für einen sekundären Server (107, 108) hat, und einen Host für einen sekundären Server betrachtet, der einen höheren Rang als ein Reserveknoten (105) hat.
- 22Appareil selon la revendication 13, dans lequel le mécanisme de sélection est configuré pour cesser (614) de communiquer des informations de rang entre le noeud donné (I) et les autres noeuds dans le système informatique distribué après que le noeud donné est disqualifié par le mécanisme de disqualification. The apparatus of claim 13, wherein the selecting mechanism is configured to cease (614) to communicate rank information between the given node (I) and the other nodes in the distributed computing system after the given node is disqualified by the disqualification mechanism. Vorrichtung nach Anspruch 13, bei der der Auswahlmechanismus so konfiguriert ist, dass er den Austausch von Ranginformationen zwischen dem gegebenen Knoten (I) und den anderen Knoten in dem verteilten Computersystem beendet (614), nachdem der gegebene Knoten durch den Ausschließungsmechanismus ausgeschlossen worden ist.
Independent claims22
55 paragraphs in 3 sections, as filed
<u style="single">Field of the Invention</u>
The present invention relates to coordinating activities between nodes in a distributed computing system. More specifically, the present invention relates to a method and an apparatus for reaching agreement between nodes in the distributed computing system regarding a node to function as a primary provider for a service.
<u style="single">Related Art</u>
As computer networks are increasingly used to link computer systems together, distributed computing systems have been developed to control interactions between computer systems. Some distributed computing systems allow client computer systems to access resources on server computer systems. For example, a client computer system may be able to access information contained in a database on a server computer system.
When a server computer system fails, it is desirable for the distributed computing system to automatically recover from this failure. Distributed computer systems possessing an ability to recover from such server failures are referred to as "highly available systems."
For a highly available system to function properly, the highly available system must be able to detect a server failure and reconfigure itself so that accesses to a failed server are redirected to a backup secondary server.
One problem in designing such a highly available system is that some distributed computing system functions must be centralized in order to operate efficiently. For example, it is desirable to centralize an arbiter that keeps track of where primary and secondary copies of a server are located in a distributed computing system. However, a node that hosts such a centralized arbiter may itself fail. Hence, it is necessary to provide a mechanism to select a new node to host the centralized arbiter.
Moreover, this selection mechanism must operate in a distributed fashion because, for the reasons stated above, no centralized mechanism is certain to continue functioning. Furthermore, it is necessary for the node selection process to operate so that the nodes that remain functioning in the distributed computing system agree on the same node to host the centralized arbiter. For efficiency reasons, it is also desirable for the node selection mechanism not to move the centralized arbiter unless it is necessary to do so.
Hence, what is needed is a method and an apparatus that operates in a distributed manner to select a node to host a primary server for a service.
Christian et al: "Fault-tolerance in air traffic control systems" ACM Transactions on Computer Systems, US, Association for Computing Machinery, New York, vol. 14, no. 3, 1 August 1996, pages 265-286, ISSN:0734-2071, describes fault tolerance in highly available distributed real time system services. Mechanisms for managing redundant server groups in a way that masks group member failures and at the same time makes the group behaviour indistinguishable from a single server are discussed. Group members are enabled to communicate with other groups to reach agreement on the state of the service they implement. The availability policy enforced for each server group will determine how many members a group should have and how closely synchronized the local states of each group member should be. A ranking system is used within a group to promote the highest-ranking server to become the primary host for the service. A Global Availability Management Service (GSAM) ensures availability of movable service groups. A Group Service Availability Manager (gSAM) ensures that for each service in the group, the prescribed availability policy is automatically enforced despite server failures and shutdowns. However, close synchronization was chosen for both the GSAM and gSAM managers, because close synchronization among the servers implementing the service yields a simpler overall design than loose synchronization, with no loss of performance. The reasons for this are that in a closely synchronized group there is no need to program the (often complicated) promotion protocols that must be executed when a server becomes the highest-ranking member after the previous leader fails. Moreover, since in a loosely synchronized approach it is the leader's role to synchronize all concurrent events affecting the group, special point-to point communication mechanisms are needed to enable ordinary group members to inform the leader about all events that they detect. This approach is said to lead to more complicated code, and also to longer delays in reconfiguring the group after component failures.
EP-A-0,750,256 (Data General Corporation) describes a framework for managing cluster membership in a multiprocessor system and for registering nodes to provide a service. In selecting nodes where registration is allowable, node preferences result from the fact that not all nodes will support all client services equally well. Node preferences may be specified as an unordered list or as an ordered list. Selection among unordered members will be influenced by recent performance characteristics of the cluster, but ordered lists are processed using a ranking system, beginning with the highest-ranking member.
SUMMARY
The present invention provides a method and system for selecting a node to host a primary server for a service from a plurality of nodes in a distributed computing system, in accordance with claims which follow. The system operates by receiving an indication that a state of the distributed computing system has changed. In response to this indication, the system determines if there is already a node hosting the primary server for the service. If not, the system selects a node to host the primary server using the assumption that a given node from the plurality of nodes in the distributed computing system hosts the primary server. The system then communicates rank information between the given node and other nodes in the distributed computing system, wherein each node in the distributed computing system has a unique rank with respect to the other nodes in the distributed computing system. The system next compares the rank of the given node with the rank of the other nodes in the distributed computing system. If one of the other nodes has a higher rank than the given node, the system disqualifies the given node from hosting the primary server.
In one embodiment of the present invention, if there exists a node to host the primary server, the system allows the node that hosts the primary server to communicate with other nodes in the distributed computing system in order to disqualify the other nodes from hosting the primary server.
In one embodiment of the present invention, the system maintains a candidate variable in the given node identifying a candidate node to host the primary server. In a variation on this embodiment, the system initially sets the candidate variable to identify the given node.
In one embodiment of the present invention, after a new node has been selected to host the primary server, if the new node is different from a previous node that hosted the primary server, the system maps connections for the service to the new node. In a variation on this embodiment, the system also configures the new node to host the primary server for the service.
In one embodiment of the present invention, the system restarts the service if the service was interrupted as a result of the change in state of the distributed computing system.
In one embodiment of the present invention, the given node in the distributed computing system can act as one of: a host for the primary server for the service; a host for a secondary server for the service, wherein the secondary server periodically receives checkpointing information from the primary server; or a spare for the primary server, wherein the spare does not receive checkpointing information from the primary server.
In one embodiment of the present invention, upon initial startup of the service, the system selects a highest ranking spare to host the primary server for the service.
In one embodiment of the present invention, the system allows the primary server to configure spares in the distributed computing system to host secondary servers for the service.
In one embodiment of the present invention, comparing the rank of the given node with the rank of the other nodes in the distributed computing system involves considering a host for a secondary server to have a higher rank than a spare.
In one embodiment of the present invention, after disqualifying the given node from hosting the primary server, the system ceases to communicate rank information between the given node and the other nodes in the distributed computing system.
BRIEF DESCRIPTION OF THE FIGURES
<ul id="ul0001" list-style="none" compact="compact"><li>FIG. 1 illustrates a distributed computing system in accordance with an embodiment of the present invention.</li><li>FIG. 2 illustrates how highly available services are controlled within a distributed computing system in accordance with an embodiment of the present invention.</li><li>FIG. 3 illustrates how a replica manager controls highly available services in accordance with an embodiment of the present invention.</li><li>FIG. 4 is a flow chart illustrating the process of selecting and configuring a new primary server in accordance with an embodiment of the present invention.</li><li>FIG. 5 is a flow chart illustrating some of the operations performed by a primary server in accordance with an embodiment of the present invention.</li><li>FIG. 6 illustrates how a node is selected to host a primary server through a disqualification process in accordance with an embodiment of the present invention.</li><li>FIG. 7 illustrates how nodes a disqualified in accordance with an embodiment of the present invention.</li></ul>
DETAILED DESCRIPTION
The following description is presented to enable any person skilled in the art to make and use the invention, and is provided in the context of a particular application and its requirements. Various modifications to the disclosed embodiments will be readily apparent to those skilled in the art, and the general principles defined herein may be applied to other embodiments and applications without departing from the scope of the present invention. Thus, the present invention is not intended to be limited to the embodiments shown, but is to be accorded the widest scope consistent with the principles and features disclosed herein.
The data structures and code described in this detailed description are typically stored on a computer readable storage medium, which may be any device or medium that can store code and/or data for use by a computer system. This includes, but is not limited to, magnetic and optical storage devices such as disk drives, magnetic tape, CDs (compact discs) and DVDs (digital versatile discs or digital video discs), and computer instruction signals embodied in a transmission medium (with or without a carrier wave upon which the signals are modulated). For example, the transmission medium may include a communications network, such as the Internet.
<u style="single">Distributed Computing System</u>
FIG. 1 illustrates a distributed computing system 100 in accordance with an embodiment of the present invention. Distributed computing system 100 includes a number of computing nodes 102-105, which are coupled together through a network 110.
Network 110 can include any type of wire or wireless communication channel capable of coupling together computing nodes. This includes, but is not limited to, a local area network, a wide area network, or a combination of networks. In one embodiment of the present invention, network 110 includes the Internet. In another embodiment of the present invention, network 110 is a local high speed network that enables distributed computing system 100 to function as a clustered computing system (hereinafter referred to as a "cluster").
Nodes 102-105 can generally include any type of computer system, including, but not limited to, a computer system based on a microprocessor, a mainframe computer, a digital signal processor, a personal organizer, a device controller, and a computational engine within an appliance.
Nodes 102-105 also host servers, which include a mechanism for servicing requests from a client for computational and/or data storage resources. More specifically, node 102 hosts primary server 106, which services requests from clients (not shown) for a service involving computational and/or data storage resources.
Nodes 103-104 host secondary servers 107-108, respectively, for the same service. These secondary servers act as backup servers for primary server 106. To this end, secondary servers 107-108 receive periodic checkpoints 120-121 from primary server 106. These periodic checkpoints enable secondary servers 107-108 to maintain consistent state with primary server 106. This makes it possible for one of secondary servers 107-108 to take over for primary server 106 if primary server 106 fails.
Node 105 can serve as a spare node to host the service provided by primary server 106. Hence, node 105 can be configured to host a secondary server with respect to a service provided by primary server 106. Alternatively, if all primary servers and secondary servers for the service fail, node 105 can be configured to host a new primary server for the service.
Also note that nodes 102-105 contain distributed selection mechanisms 132-135, respectively. Distributed selection mechanisms 132-135 communicate with each other to select a new node to host primary server 106 when node 102 fails or otherwise becomes unavailable. This process is described in more detail below with reference to FIGs. 2-6.
<u style="single">Controlling Highly Available Services</u>
FIG. 2 illustrates how highly available services 202-205 are controlled within distributed computing system 100 in accordance with an embodiment of the present invention. Note that highly available services 202-205 continue to operate even if individual nodes of distributed computing system 100 fail.
Highly available services 202-205 operate under control of replica manager 206. Referring to FIG. 3, for each service, replica manager 206 keeps a record of which nodes in distributed computing system 100 function as primary servers, and which nodes function as secondary servers. For example, in FIG. 3 replica manager 206 keeps track of highly available services 202-205. The primary server for service 202 is node 103, and the secondary servers are nodes 104 and 105. The primary server for service 203 is node 104, and the secondary servers are nodes 103 and 105. The primary server 106 for service 204 is node 102, and the secondary servers 107,108 are nodes 103 and 104. The primary server for service 205 is node 103, and the secondary servers are nodes 102, 104 and 105.
Replica manager 206 additionally performs a number of related functions, such as configuring a node to host a primary (which may involve demoting a current host for the primary to host a secondary). Replica manager 206 may additionally perform other functions, such as: adding a service; removing providers for a service; registering providers for a service; removing a service; handling provider failures; bringing up new providers for a service; and handling dependencies between services (which may involve ensuring that primaries for dependent services are co-located on the same node).
Referring back to FIG. 2, replica manager 206 is itself a highly available service that operates under control of replica manager manager (RMM) 208. Note that RMM 208 is not managed by a higher level service.
As illustrated in FIG. 2, RMM 208 communicates with cluster membership monitor (CMM) 210. CMM 210 monitors cluster membership and alerts RMM 208 if any changes in the cluster membership occur.
CMM 210 uses transport layer 212 to exchange messages between nodes 102-105.
<u style="single">Process of Selecting a New Primary</u>
FIG. 4 is a flow chart illustrating the process of selecting and configuring a primary server in accordance with an embodiment of the present invention. Note that this process is run concurrently by each active node in distributed computing system 100. The system begins by receiving an indication from CMM 210 that the membership in the cluster has changed (step 401).
In response to this indication, the system obtains a lock on a local candidate variable which contains an identifier for a candidate node to host primary server 106 (step 402). The system also obtains an additional lock to hold off requesters for the service (step 404).
Next, the system executes a disqualification process by communicating with other nodes in distributed computing system 100 in order to disqualify the other nodes from acting as the primary server 106 (step 406). This process is described in more detail with reference to FIG. 6 below.
After the disqualification process, the remaining node, which is not disqualified, becomes the primary node. If the node hosting primary server 106 has changed, this may involve re-mapping connections for the service to point to the new node (step 408). It may also involve initializing the new node to act as the host for the primary (step 410).
Finally, the service is started (step 412). This may involve unfreezing the service if it was previously frozen, as well as releasing the previously obtained lock that holds off requesters for the service.
FIG. 5 is a flow chart illustrating some of the operations performed by a primary server 106 in accordance with an embodiment of the present invention. During operation, primary server 106 performs periodic checkpointing operations 120-121 with secondary servers 107-108, respectively (step 502). These checkpointing operations allow secondary servers 107-108 to take over from primary server 106 if primary server 106 fails. Primary server 106 also periodically attempts to promote spare nodes (such as node 105 in FIG. 1) to host secondaries (step 504). This promotion process involves transferring state information to a spare node in order to bring the spare node up to date with respect to the existing secondaries.
FIG. 6 illustrates how a node is selected to host primary server 106 through a disqualification process in accordance with an embodiment of the present invention. Note that FIG. 6 describes in more detail the process described above with reference to step 406 in FIG. 4.
The system starts by determining if a node that was previously hosting primary server 106 continues to exist (step 602).
If not, the system retrieves the state of a local provider for the service (step 604). The system then sets the candidate variable to identify the local provider (step 606), and subsequently unlocks the candidate lock that was set previously in step 402 of FIG. 4 (step 608). Next, if the candidate is not the local provider, the system ends the process (step 610).
Next, for all other nodes I in the cluster, the system attempts to disqualify node I by writing a new identifier into the candidate variable for node I if the rank of node I is less than the rank of the present node (step 612). This process is described in more detail with reference to FIG. 7 below. Finally, if node I's local provider has a higher rank than the present node, the process terminates because the present node is disqualified (step 614).
Note that a rank of a node can be obtained by comparing a unique identifier for the node with unique identifiers for other nodes. Also note that the rank of a primary server is greater than the rank of a secondary server, and that the rank of a secondary server is greater than the rank of a spare. The above-listed restrictions on rank ensure that an existing primary that has not failed continues to function as the primary, and that an existing secondary will be chosen ahead of a spare. Of course, when the system is initialized, no primaries or secondaries exist, so a spare is selected to be the primary.
On the other hand, if the node that was hosting primary server 106 continues to function, the system sets the candidate to be this node (step 616), and unlocks the candidate node (step 618).
If the present node does not host primary server 106, the process is finished. Otherwise, if the present node is hosting primary server 106, the system considers each other node I in the cluster. If the present node has already communicated with I, the system skips node I (step 622). Otherwise, the system communicates with node I in order to disqualify node I from acting as the host for primary server 106 (step 624). This may involve causing an identifier for the present node to be written into the candidate variable for node I.
FIG. 7 illustrates how nodes are disqualified in accordance with an embodiment of the present invention. Note that FIG. 7 describes in more detail the process described above with reference to step 612 in FIG. 6. The caller first locks the candidate variable for node I (step 702). If the caller determines that the caller's provider has a higher rank than is specified in the candidate variable for I, the caller overwrites the candidate variable for I with the caller's provider (step 704). Next, the caller unlocks the candidate variable for I (step 706).
The foregoing descriptions of embodiments of the invention have been presented for purposes of illustration and description only. They are not intended to be exhaustive or to limit the present invention to the forms disclosed. Accordingly, many modifications and variations will be apparent to practitioners skilled in the art. Additionally, the above disclosure is not intended to limit the present invention. The scope of the present invention is defined by the appended claims.
Contents3
5 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5
Every citation, both waysCites: the store holds 4 of 5
| Document | Relation | Office |
|---|---|---|
| EP0750256A | Cites | European Patent Office (EPO) |
| WO9749039A | Cites | World Intellectual Property Organization (WIPO) |
| US5526492A | Cites | United States of America |
| US5793968A | Cites | United States of America |
7 members in 4 offices
Priority claims10
| Document | Office | Kind | Date |
|---|---|---|---|
| 160992P | United States of America | – | |
| 16099299 | United States of America | P | |
| 16099299 | United States of America | P | |
| 662553 | United States of America | – | |
| 66255300 | United States of America | A | |
| 66255300 | United States of America | A | |
| 160992P | – | – | – |
| 662553 | – | – | – |
| US19990160992P | – | – | – |
| US20000662553 | – | – | – |
Members7
| Document | Office | Kind | |
|---|---|---|---|
| EP1096751A2 | European Patent Office (EPO) | A2 | |
| EP1096751A3 | European Patent Office (EPO) | A3 | |
| US6957254B1 | United States of America | B1 | |
| EP1096751B1This record | European Patent Office (EPO) | B1 | |
| AT347224T | Austria | T | |
| ATE347224T1 | Austria | T1 | |
| DE60032099D1 | Germany | D1 |
48 legal events, as 4 offices reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | Office | |
|---|---|---|---|
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Patent expired after termination of 20 yearsExpiredPE20 | PE20 | GB | |
| Annual fee paid to national office [announced via postgrant information from national office to epo]GrantedPGFP | PGFP | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| No opposition filedOpposition26N | 26N | EP | |
| No opposition filed within time limitOppositionORIGINAL CODE: 0009261PLBE | PLBE | EP | |
| Information on the status of an ep patent application or granted ep patentGrantedSTATUS: NO OPPOSITION FILED WITHIN TIME LIMITSTAA | STAA | EP | |
| Fr: translation not filedEN | EN | EP | |
| Patent ceasedCeasedPL | PL | CH | |
| Nl: lapsed or annulled due to failure to fulfill the requirements of art. 29p and 29m of the patents actLapsedNLV1 | NLV1 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Corresponds to:REF | REF | EP | |
| European patents granted designating irelandGrantedFG4D | FG4D | IE | |
| European patent takes effect as a national patent in ch/liEP | EP | CH | |
| European patents granted designating irelandGrantedFG4D | FG4D | IE | |
| Designated contracting statesAK | AK | EP | |
| European patent grantedGrantedFG4D | FG4D | GB | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| Lapsed in a contracting state [announced via postgrant information from national office to epo]LapsedPG25 | PG25 | EP | |
| (expected) grantORIGINAL CODE: 0009210GRAA | GRAA | EP | |
| Grant fee paidORIGINAL CODE: EPIDOSNIGR3GRAS | GRAS | EP | |
| Despatch of communication of intention to grant a patentORIGINAL CODE: EPIDOSNIGR1GRAP | GRAP | EP | |
| Party data changed (applicant data changed or rights of an application transferred)RAP1 | RAP1 | EP | |
| First examination report despatched17Q | 17Q | EP | |
| Designation fees paidAT BE CH CY DE DK ES FI FR GB GR IE IT LI LU MC NL PT SEAKX | AKX | EP | |
| Request for examination filed17P | 17P | EP | |
| Designated contracting statesAK | AK | EP | |
| Request for extension of the european patentAL;LT;LV;MK;RO;SIAX | AX | EP | |
| Information provided on ipc code assigned before grant7H 04L 29/06 A, 7G 06F 9/50 B, 7G 06F 11/20 BRIC1 | RIC1 | EP | |
| Search report despatchedORIGINAL CODE: 0009013PUAL | PUAL | EP | |
| Designated contracting statesAK | AK | EP | |
| Request for extension of the european patentAL;LT;LV;MK;RO;SIAX | AX | EP | |
| Public reference made under article 153(3) epc to a published international application that has entered the european phaseORIGINAL CODE: 0009012PUAI | PUAI | EP |
Numbers
- Publication
- 1096751
- Publication, DOCDB
- 1096751
- Publication, EPODOC
- EP1096751
- Application
- 203642
- Application, DOCDB
- 00203642
- Application, EPODOC
- EP20000203642
Titles3
- German
- Verfahren und Vorrichtung zum Erreichen einer Übereinkunft zwischen Knoten in einem verteilten System
- English
- Method and apparatus for reaching agreement between nodes in a distributed system
- French
- Procédé et dispositif pour atteindre un accord entre noeuds dans un système distribué
Classification
- CPC, 7
- G06F9/5061
- H04L41/0856
- H04L41/5054
- H04L67/1002
- H04L67/1034
- H04L69/40
- H04L67/1001
- IPC, 4
- H04L29 06
- G06F9 50
- G06F11 20
- H04L69 40
Designated states19
- Contracting states, 19
- Austria
- Belgium
- Switzerland
- Cyprus
- Germany
- Denmark
- Spain
- Finland
- France
- United Kingdom
- Greece
- Ireland
- Italy
- Liechtenstein
- Luxembourg
- Monaco
- Netherlands (Kingdom of the)
- Portugal
- Sweden
