Induction of a node into a group
20 claims: 8 independent, 12 dependent
- 1誘導者ノードが新入者ノードを分散型コンピューターシステムに導入するためにコンピューターで実施される方法であって、 前記分散型コンピューターシステムへの前記新入者ノードの導入を完了するために必要な情報を含む導入タスクを生成すること、 コンピューターネットワークを介して、前記新入者ノードに前記導入タスクを送ること、 前記新入者ノードを識別する情報、及び、前記新入者ノードとの通信を可能にするに足る情報を含むメンバーシップ要求を、前記新入者ノードから前記ネットワークを介して受け取ること、 誘導者ノード及び新入者ノードの役割を定義するブートストラップメンバーシップを生成し、前記ブートストラップメンバーシップを 配置 し、前記新入者ノードに対してブートストラップメンバーシップレディメッセージを送ること、 生成された前記ブートストラップメンバーシップを参照する決定論的な状態マシンを生成すること、及び 前記新入者ノードが対応するブートストラップメンバーシップを生成したというアクナレッジメントを受信すること を備える方法。
- 2前記導入タスクの 生 成は 、 前記導入タスクが 少なくとも1つの導入チケットを含む ように 実行される 請求項1に記載のコンピューターで実施される方法。
- 3前記導入チケットは、暗号化されたファイルとして構成される 請求項 2 に記載のコンピューターで実施される方法。
- 4前記導入タスクは、少なくとも1つの導入後タスクを含む 請求項2に記載のコンピューターで実施される方法。
- 5前記新入者ノードから受信された前記メンバーシップ要求を認証すること、及び 前記メンバーシップ要求が無効であったら、前記導入を終了させること をさらに備える請求項1に記載のコンピューターで実施される方法。
- 6前記ブートストラップメンバーシップはさらに、 決定論的に生成されたメンバーシップ識別、及び 前記誘導者ノード及び前記新入者ノードの役割を含む 請求項1に記載のコンピューターで実施される方法。
- 7前記新入者ノードが気付くべき少なくとも1つのノード及びロケーションを前記新入者ノードに対して送ること をさらに備える請求項1に記載のコンピューターで実施される方法。
- 8前記導入タス クは 、少なくとも前記誘導者ノードの再起動の間持続す る 請 求項1に記載のコンピューターで実施される方法。
- 9メモリと、 プロセッサとを備える計算装置であって、 前記プロセッサは、分散型コンピューターシステムに新入者ノードを導入するよう構成された誘導者ノードとして前記計算装置を実行するために前記メモリ内に記憶される命令を実行するよう構成され、 記憶される前記命令は、前記プロセッサに 前記分散型コンピューターシステムへの前記新入者ノードの導入を完了するために必要な情報を含む導入タスクを生成すること、 コンピューターネットワークを介して、前記新入者ノードに前記導入タスクを送ること、 前記新入者ノードを識別する情報、及び、前記新入者ノードとの通信を可能にするに足る情報を含むメンバーシップ要求を、前記新入者ノードから前記ネットワークを介して受け取ること、 誘導者ノード及び新入者ノードの役割を定義するブートストラップメンバーシップを生成し、前記ブートストラップメンバーシップを 配置 し、前記新入者ノードに対してブートストラップメンバーシップレディメッセージを送ること、 生成された前記ブートストラップメンバーシップを参照する決定論的な状態マシンを生成すること、及び 前記新入者ノードが対応するブートストラップメンバーシップを生成したというアクナレッジメントを受信すること を実行させるよう構成される計算装置。
- 10分散型コンピューターシステムに新入者ノードを導入するよう構成された誘導者ノードとして計算装置を構成するデータ及び命令を記憶する有形のデータ記憶メディアであって、 記憶される前記データ及び命令は、前記計算装置に、 前記分散型コンピューターシステムへの前記新入者ノードの導入を完了するために必要な情報を含む導入タスクを生成すること、 コンピューターネットワークを介して、前記新入者ノードに前記導入タスクを送ること、 前記新入者ノードを識別する情報、及び、前記新入者ノードとの通信を可能にするに足る情報を含むメンバーシップ要求を、前記新入者ノードから前記ネットワークを介して受け取ること、 誘導者ノード及び新入者ノードの役割を定義するブートストラップメンバーシップを生成し、前記ブートストラップメンバーシップを 配置 し、前記新入者ノードに対してブートストラップメンバーシップレディメッセージを送ること、 生成された前記ブートストラップメンバーシップを参照する決定論的な状態マシンを生成すること、及び 前記新入者ノードが対応するブートストラップメンバーシップを生成したというアクナレッジメントを受信すること を実行させるよう構成されるデータ記憶メディア。
- 11新入者ノードが誘導者ノードによって分散型コンピューターシステム内に導入されるためにコンピューターで実施される方法であって、 導入チケットを少なくとも含む導入タスクの受領を待機すること、 コンピューターネットワークを介して前記導入タスク及び前記導入チケットを受け取り、前記導入チケット内の情報を用いてアプリケーションプラットホームを設定すること、 前記誘導者ノードに対して前記新入者ノードが導入を開始していることを通知するように構成されたブートストラップメンバーシップ要求を、前記 誘導者 ノードに対し、前記コンピューターネットワークを介して送ること、 前記誘導者ノードからブートストラップメンバーシップレディメッセージを受け取り、ブートストラップメンバーシップを生成及び 配置 し、前記ブートストラップメンバーシップを参照する決定論的な状態マシンを生成するとともに、前記誘導者ノードに対して 前記 ブートストラップメンバーシップを 生成したというアクナレッジメント を知らせること を備える方法。
- 12送信することは、受領された前記導入タスクの識別、及び、前記新入者ノードとの通信を可能にする情報を含む前記ブートストラップメンバーシップ要求を用いて行われる 請求項11に記載のコンピューターで実施される方法。
- 13前記決定論的な状態マシンが、前記誘導者ノードから、前記新入者ノードが気付くべきロケーション及びノードのリストを受領すること をさらに備える請求項11に記載のコンピューターで実施される方法。
- 14前記決定論的な状態マシンが、前記誘導者ノードから、前記新入者ノードが気付くべきロケーション及びノードのリストを含む提案を受領すること をさらに備える請求項11に記載のコンピューターで実施される方法。
- 15受領された前記導入チケットを認証すること をさらに備える請求項11に記載のコンピューターで実施される方法。
- 16前記ブートストラップメンバーシップ要求は、前記導入タスクの識別、ノード、ロケーション、前記新入者ノードのホストネーム及びポートのうちの少なくとも1つを含む 請求項11に記載のコンピューターで実施される方法。
- 17前記ブートストラップメンバーシップを生成することは、決定論的に生成されたメンバーシップ識別と、前記誘導者ノード及び前記新入者ノードの役割とを用いて前記ブートストラップメンバーシップを生成することを含む 請求項11に記載のコンピューターで実施される方法。
- 18受領された前記導入タスク内に特定されるタスクを実行すること をさらに備える請求項11に記載のコンピューターで実施される方法。
- 19メモリと、 プロセッサとを備える計算装置であって、 前記プロセッサは、 誘導者 ノードによって分散型コンピューターシステム内に導入される新入者ノードとして前記計算装置を構成するために前記メモリ内に記憶される命令を実行するよう構成され、 記憶される前記命令は、前記プロセッサに 導入チケットを含む導入タスクの受領を待機すること、 コンピューターネットワークを介して前記導入タスク及び前記導入チケットを受け取り、前記導入チケット内の情報を用いてアプリケーションプラットホームを設定すること、 前記誘導者ノードに対して前記新入者ノードが導入を開始していることを通知するように構成されたブートストラップメンバーシップ要求を、前記 誘導者 ノードに対し、前記コンピューターネットワークを介して送ること、 前記誘導者ノードからブートストラップメンバーシップレディメッセージを受け取り、ブートストラップメンバーシップを生成及び 配置 し、前記ブートストラップメンバーシップを参照する決定論的な状態マシンを生成するとともに、前記誘導者ノードに対して前記ブートストラップメンバーシップを 生成したというアクナレッジメント を知らせること を実行させるよう構成される計算装置。
- 20誘導者ノードによって分散型コンピューターシステム内に導入される新入者ノードとして計算装置を構成するデータ及び命令を記憶する有形のデータ記憶メディアであって、 記憶される前記命令は、前記計算装置に、 導入チケットを含む導入タスクの受領を待機すること、 コンピューターネットワークを介して前記導入タスク及び前記導入チケットを受け取り、前記導入チケット内の情報を用いてアプリケーションプラットホームを設定すること、 前記誘導者ノードに対して前記新入者ノードが導入を開始していることを通知するように構成されたブートストラップメンバーシップ要求を、前記 誘導者 ノードに対し、前記コンピューターネットワークを介して送ること、 前記誘導者ノードからブートストラップメンバーシップレディメッセージを受け取り、ブートストラップメンバーシップを生成及び 配置 し、前記ブートストラップメンバーシップを参照する決定論的な状態マシンを生成するとともに、前記誘導者ノードに対して前記ブートストラップメンバーシップを 生成したというアクナレッジメント を知らせること を実行させるよう構成されるデータ記憶メディア。
Independent claims20
120 paragraphs, as filed
0001This application claims the interests of US Provisional Application No. 61 / 746,867 filed December 28, 2012 and US Application No. 13 / 835,888 filed March 15, 2013.
0002Collaborative projects, which are often carried out in parallel among multiple globally separated resources (ie, multi-site collaborative projects), have become commonplace in various types of projects. Examples of such projects include, but are not limited to, software development, jet airliner design, and automobile design. Reliance on distributed resources to accelerate the project timeline through optimizing human resource utilization and leveraging global resource skill sets has proven to provide favorable results.
0003The decentralized computer solution used when proceeding with a multi-site joint project will be referred to as a decentralized multi-site joint computer solution here. However, the decentralized multi-site collaborative computer solution is just one example of a decentralized computer solution. In one example, a decentralized computer solution includes a network of computers operating a car. In another example, a distributed computer solution includes a network of computers in one geographic location (data center). In yet another example, a distributed computer solution is multiple computers (ie, subnets) connected to a router.
0004Although decentralized computer solutions have traditionally existed, they cannot be without restrictions that adversely affect their effectiveness, reliability, usefulness, scalability, transparency, and / or security. In particular, traditional decentralized multi-site collaborative computer solutions have limited ability to synchronize work from globally distributed development bases in real time and to fault tolerance. This limitation often forces changes in software development and delivery procedures that cause delays and increase risk. Therefore, the cost savings and productivity gains that would be achieved by running collaborative projects with traditional distributed computer solutions are not fully realized.
0005Traditional distributed multi-site collaborative computer solutions undesirably force users to change development procedures. For example, a traditional distributed multi-site collaborative computer solution that lacks the benefits associated with real-time information management capabilities has local and remote Concurrent Versions System (CVS) repositories synchronized at any given time. It has the fundamental problem of not being able to guarantee that. This means that developers from different sites are more likely to inadvertently overwrite and collide with each other's work. To prevent the risk of such overwrites and conflicts, traditional distributed multi-site collaborative computer solutions include excessive and / or error-prone source code branching and manual file merging as part of the development process. Needed as a department. This effectively forces the development business to be split based on time zone, making coordination between distributed development teams very difficult, if not impossible.
0006The replication state machine is a preferred implementation of a distributed computer solution. Also, one of the few possible examples of distributed computer solutions is the replication information repository. Therefore, the replication state machine, more specifically, is a preferred implementation of the replication information repository. Also, one of several possible applications for replication information repositories is a distributed multi-site collaborative computer solution. Therefore, the replication state machine, more specifically, is a preferred implementation of a distributed multi-site joint computer solution.
0007Therefore, distributed computer solutions often rely on replication state machines, replication information repositories, or both. The replication state machine and / or replication information repository provide concurrent information generation, manipulation, and management, and are therefore an important aspect of most distributed computer solutions. However, known techniques for facilitating state machine replication as well as information repository replication are not without flaws.
0008Traditional implementations that facilitate replication of state machines have one or more flaws that limit their effectiveness. One such flaw is that the consensus protocol tends to have repeated preemptions on the proposing side, which adversely affects extensibility. Another such flaw is weak leader optimization. The execution of optimization) requires the election of leaders, which causes such optimizations to adversely affect complexity, speed, and scalability, and one or more messages per consent. It is necessary (eg, 4 instead of 3), which adversely affects speed and scalability. Another such flaw is that agreements must be reached in sequence, which negatively impacts speed and scalability. Another such flaw is that there is a limit to the reuse of fixed storage, and if any of these flaws are present, this imposes a heavy burden on deployment. This is because the storage needs for such arrangements will grow continuously, potentially, and endlessly. Another such flaw is that there is a limit to the effective handling of large proposals and many small proposals, and if any of these flaws are present, they adversely affect scalability. Another such flaw is that a relatively large number of messages must be communicated to facilitate state machine replication, which adversely affects scalability and wide area network compatibility. Another limitation is that communication message delays adversely affect scalability. Another such flaw is that addressing failure scenarios by dynamically changing participants (eg, joining and excluding as needed) on replicated machines adversely affect complexity and scalability. It means to exert.
0009Traditional implementations that facilitate the replication of information repositories have one or more flaws that limit their effectiveness. One such flaw is that some traditional multi-site collaborative computer solutions require a single central coordinator to facilitate the replication of centrally coordinated information repositories. The central coordinator adversely affects scalability in an undesired way, as all updates to the information repository must be delivered via a single central coordinator. Moreover, such implementations are not very useful because the failure of a single central coordinator would stop the update of any replica of the information repository. Another such flaw is that in performing information repository replication that relies on log regeneration, information repository replication is facilitated in an active-passive manner. Therefore, only one copy can be updated at any time. As a result, other replicas can only be idle or provide read-only applications, such as data mining applications, resulting in poor resource utilization. The other such flaw is that the implementation is weakly consistent, backed by conflict-resolution heuristics and / or application-intervention mechanisms. Occurs when relying on replication). This type of information repository allows conflicting updates for replication of the information repository and requires applications that utilize the information repository to resolve these conflicts. Therefore, such an implementation adversely affects transparency to the application.
0010Further referring to the fact that traditional implementations that facilitate replication of information repositories have one or more flaws that limit their effectiveness, implementations that rely on disk mirroring solutions are known to have one or more flaws. ing. This type of implementation is an active-passive implementation. Therefore, one such flaw is that only one copy can be used by the application at any given time. As a result, other replicas (ie, passive mirrors) cannot be read or written while acting as passive mirrors, resulting in poor resource utilization. Another flaw in this special implementation is that the replication method is unaware of the transaction limits of the application. Therefore, at the time of failure, the mirror has a partial result of the transaction and may therefore be unavailable. Another such flaw is that the replication method propagates the change from the node where the change to the information occurred to all other nodes. Such implementations require unnecessarily large bandwidth, as the size of changes to information is often significantly larger than the size of the command that caused the change. Another such flaw is that if the information in the master repository is corrupted for any reason, the corruption will be propagated to all replicas of that repository. As a result, the information repository becomes unrecoverable or must be recovered from an old backup copy, thus causing further information loss.
0011Therefore, a replication state machine that overcomes the drawbacks associated with conventional replication state machines is useful and advantageous. More specifically, the replication information repository constructed using such a replication state machine is superior to the conventional replication information repository. More specifically, replicated CVS repositories built using such replicated state machines are superior to traditional replicated CVS repositories.
0012The use of distributed computer solutions as described above therefore provides a relatively effective and efficient means of sharing information between physically separated locations, logically separated locations, and so on. In terms of providing, it was the key to the success of such a joint project. Each such location can have one or more computer nodes in a distributed computer system. New nodes wishing to participate in a joint project need to be invited to join an existing set of nodes and can be seen from there, and the newly invited nodes can exchange messages and interact with each other. Need to be informed about the locations and nodes that are allowed to act on.
0013<figref num="1">It is a block diagram which shows the functional relationship of the element in the multi-site computer system architecture based on one Embodiment.</figref>
0014<figref num="2">It is a high-level block diagram which shows the arrangement of the elements which make up a multi-site computer system architecture based on one Embodiment.</figref>
0015<figref num="3">It is a block diagram which shows the functional component of the replication state machine based on one Embodiment.</figref>
0016<figref num="4">It is a block diagram which shows the proposal issued by the local application node by one Embodiment.</figref>
0017<figref num="5">It is a block diagram which shows the entry structure of the global sequencer of the replication state machine shown in FIG.</figref>
0018<figref num="6">It is a block diagram which shows the entry structure of the local sequencer of the replication state machine shown in FIG.</figref>
0019<figref num="7">It is a block diagram which shows the replicator based on one Embodiment.</figref>
0020<figref num="8">It is a detail level block diagram which shows the arrangement of the elements which make up a multi-site computer system architecture based on one Embodiment.</figref>
0021<figref num="9">It is a figure which shows one aspect of the apparatus, method, and system which realizes safe and authenticated introduction of a node into a group of nodes by one Embodiment.</figref>
0022<figref num="10">At the same time, it is a block diagram of a computing device in which an embodiment can be implemented.</figref>
0023Disclosed herein are various aspects that facilitate the practical implementation of replica state machines in various distributed computer system architectures (eg, distributed multisite co-computer system architectures). Those skilled in the art will be aware of one or more traditional implementations of replica state machines. Traditional implementations of such state machines are, for example, FB Schneider (FB). Written by Schneider) and published at ACM Computing Surveys 22 in December 1990, "Implementing fault-tolerant services using the It is disclosed in a publication entitled "state machine approach: A tutorial" (pages 299-319). The contents of this publication are incorporated herein by reference as they form part of this specification. This embodiment enhances the scalability, reliability, availability, and fault tolerance aspects of the traditional implementation of state machines in a distributed application system architecture, as discussed in more detail below. ..
0024This embodiment provides a practical implementation of a replica state machine within various distributed computer system architectures (eg, distributed multi-site co-computer system architecture). More specifically, this embodiment enhances the scalability, reliability, availability, and fault tolerance of replication state machines and / or replication information repositories in a distributed computer system architecture. Therefore, this embodiment advantageously overcomes one or more shortcomings associated with conventional approaches for running replica state machines and / or replica information repositories in distributed computer system architectures.
0025In one embodiment, the replication state machine can include a proposal manager, a consent manager, a collision / backoff timer, and a storage reclaimer. The Proposal Manager facilitates the management of proposals issued by a node in a distributed application so that the proposals are executed cooperatively by all nodes in the distributed application that need to execute the proposal. .. All nodes probably include themselves, but they are not required. The consent manager facilitates consent to the proposal. The collision / backoff timer eliminates repeated round preemptions in an attempt to achieve consent to the proposal. The storage reclaimer reuses the fixed storage used to accept the proposal and / or store the proposal.
0026In other embodiments, the decentralized computer system architecture can include a network system and a plurality of decentralized computer systems interconnected via the network system. Each of the distributed computer systems can include each replication state machine and each local application node connected to each replication state machine. Multiple distributed computer systems Each replication state machine facilitates the management of proposals and agrees to proposals to allow the distributed application nodes of all other distributed computer systems to execute proposals in a coordinated manner. Promote, eliminate repeated round preemptions in attempts to achieve consent to a proposal, and reuse fixed storage used to store at least one of the proposal consent and proposals.
0027In other embodiments, the method comprises a plurality of operations. An operation may be performed to facilitate agreement on proposals received from the local application node. An operation may be performed to eliminate repeated round preemptions in an attempt to achieve consent to the proposal. An operation may be performed to reuse the fixed storage used to store the consent of the proposal and at least one of the proposals.
0028In at least one embodiment, at least a portion of the proposal comprises the proposed steps corresponding to the execution of information updates invoked by the nodes of the distributed application. The order of issuance of proposals may be preserved while concurrent consent to the proposals is in progress. Part of the proposal can be a proposed write step corresponding to each information update, and the proposal manager assigns a local serial number to each of those proposed write steps and distributes to perform the proposed write step. Generate a globally unique interleaving of the proposed write steps so that all nodes of the type application can perform the proposed write steps in a common sequence. There may be a local sequencer containing multiple entries, each associated with a corresponding one in multiple proposals, as well as a global sequencer containing multiple entries referencing the corresponding one in each local sequencer entry. obtain. Each of the entries in the local sequencer may have a unique local sequence number assigned to it, or each entry in the local sequencer may be sequentially arranged with respect to the sequence number assigned to it. Then, after the consent manager has proceeded with the consent of one of the proposals, the entry corresponding to the proposed proposal determines the position in the global sequencer where the entry is located. Can be generated within. The storage reclaimer can reuse the fixed storage by locating the entry in the global sequencer, making it known to all nodes, and then deleting the record for the proposal from the fixed proposal storage. The collision / backoff timer elapses after the current one of the rounds for the first proposer has started and before the next one of the rounds for the first proposer is started. An operation that waits for the duration of the preemption delay calculated as the time to do
0029Then referring to the drawings, FIG. 1 shows a multi-site computer system architecture based on one embodiment (ie, referred to herein as the multi-site computer system architecture 100) with each other by a wide area network (WAN) 110. It shows that it can include multiple distributed application systems 105 connected. Each of the plurality of distributed application systems 105 may include a plurality of distributed application nodes 115 (eg, an application running on a workstation), a replicator 120, and a repository replica 125. The replicator 120 of each distributed application system 105 may be connected between the WAN 110, the distributed application node 115 of each distributed application system 105, and the repository replica 125 of each distributed application system 105.
0030In one embodiment, each repository replica 125 is a Concurrent Versions System (CVS) repository. CVS is well-known open source code that manages system versions. CVS acts as a central server to which multiple CVS clients (eg, distributed application nodes 115) connect using, for example, the CVS protocol over Communication Control Protocol (TCP), as well as source code for versioning many other systems. Designed to do. When executed, the CVS server separates processes for each client connection and handles CVS requests from each client. Therefore, Replicator 120 and Repository Replica 125 allow multiple replicas of the CVS repository. While the CVS information repository is an example of an information repository useful in one embodiment, the subject matter of this disclosure is useful in replicating other types of information repositories. Databases and file systems are examples of such other types of information repositories. Therefore, the usefulness and applicability of the embodiments are not limited to a particular type of information repository.
0031As discussed in detail below, each replicator 120 may be configured to write information updates from each distributed application system 105 to each repository replica 125 of another distributed application system 105. Each replicator 120 may be an intermediary acting as an application gateway between the CVS client (ie, each distributed application system 105) and the default CVS server (ie, each repository replica 125). Each Replicator 120 works with other equivalent Replicators to ensure that all Repository Replicas 125 remain in sync with each other.
0032Unlike traditional solutions, Multisite Computer System Architecture 100 does not rely on a central transaction coordinator known as a single point of failure. The Multi-Site Computer System Architecture 100 provides a unique approach to real-time and active-active replication, the principle of one replication equivalent across all CVS repository replicas of a distributed application system. Works on the basis of. Therefore, based on one embodiment, all repository replicas synchronize in real time with all other repository replicas and are therefore based on information from users on all nodes of the distributed application system (eg, from the same code base). Programmers who work with) always work based on information from the same information base.
0033Through the integration of Replicator 120 and each Repository Replica 125, each Repository Replica becomes an active node on WAN 110 with its own transaction coordinator (ie, each Replicator 120). Each distributed transaction coordinator accepts local updates and communicates them in real time to all other repository replicas 125. Therefore, all users within the Multisite Computer System Architecture 100 work efficiently based on the same repository information (ie, a single CVS information repository), regardless of location. To that end, a multi-site computer system architecture based on one embodiment is a cost-effective, fault-tolerant software configuration management (SCM) solution that synchronizes the work of globally distributed development teams in real time. ..
0034In the event of a network or server failure, developers can continue to work. Changes are recorded in the local Replicator 120 transaction journal. Transaction journals are functionally similar to database redo logs. When connectivity is restored, the local Replicator 120 contacts the Replicator 120 of another distributed application system 105, keeping itself up-to-date and local while the network or system is down. Apply the changes captured in the transaction journal. Recovery can be performed automatically without the intervention of a CVS administrator. This self-healing capability ensures zero data loss and no loss of development time, eliminating the risk of human error in disaster recovery scenarios.
0035The benefit of working based on essentially the same repository information is that you don't have to change development procedures when development moves abroad, and big builds are completed when work from multiple sites is merged. This includes not having to play while waiting to do so, early detection of development problems, and less resources devoted to quality assurance (eg, reducing the use of redundant resources). In addition, disaster recovery is not a problem as the integrated self-healing capability provides failure avoidance. No work is lost when the system goes down.
0036As disclosed above, the implementation of a replica state machine based on one embodiment has a favorable impact on the scalability, reliability, usefulness, and fault tolerance of such a replica state machine. By positively impacting scalability, reliability, usefulness, and fault tolerance, this embodiment provides a practical approach for implementing replica state machines within a multi-site computer system architecture. An implementation of a replica state machine according to an embodiment may satisfy all or part of the following objectives: so that the nodes of a distributed computer system consisting of multiple computers can work together to evolve their state. To ensure that the integrity of a distributed system of multiple computers is preserved, regardless of any or partial failure of the computer network, computer, or computer resources, consists of distributed application nodes. Allowing a reliable system to generate from a component with unobtrusive reliability, guaranteeing the termination of the consent protocol with a probability as a function of the time it takes to approach 1 gradually, regardless of the conflict of the consent protocol. That, eradicating consent protocol conflicts under normal operating conditions, improving the efficiency of consent protocols, reducing and suppressing the memory and disk usage of replication state machines, and networking by replication state machines. To reduce resource usage, increase the throughput of state transitions achievable by replication state machines, and enable more efficient management of memory and disk resources by distributed applications dominated by replication state machines.
0037As shown in FIG. 2, the multi-site computer function based on one embodiment is facilitated by a plurality of replication state machines 200 that exchange information with each other and with each local application node 205 through the network system 210. Ru. Preferably, but not required, each local application node 205 is that of a distributed application and may function as a proposer or recipient of a proposal at any time in time. In one embodiment, the network system 210 is a wide area network (WAN) connected between replication state machines 200 and a local area network connected between each replication state machine 200 and each local application node 205. (LAN) and can be included. For example, each replication state machine 200 and each local application node 205 is located at each site for one multi-site collaborative computer project. The LAN portion of network system 210 facilitates the sharing of information by local standards (ie, between each replication state machine 200 and each local application node 205), and the WAN portion of network system 210 is global standards (ie,). , Promote information sharing between 200 replication state machines). While LAN, WAN, or both are components of a network system based on one embodiment, the embodiments are not limited to any particular form of the network. For example, another embodiment of a network system based on one embodiment is an ad hoc network system that includes a computer embedded in a car, a network system that includes multiple subnets in a data center, or one subnet in a data center. Includes network systems including.
0038FIG. 3 is a block diagram showing the functional components of each replication state machine 200 shown in FIG. Each replication state machine 200 has Proposal Manager 220, Fixed Proposal Storage 230, Consent Manager 240, Consent Store 245, Distributed File Transfer Protocol (DFTP) Layer 250, Collision & Backoff Timer 260, Local Sequencer 270, Global Sequencer 280, And storage reclaimer 290 (ie, fixed storage garbage collector). Proposal Manager 220, Fixed Proposal Storage 230, Consent Manager 240, Consent Store 245, Distributed File Transfer Protocol (DFTP) Layer 250, Collision & Backoff Timer 260, Local Sequencer 270, Global Sequencer 280, and Storage Reclaimer 290 They are interconnected at least in part of each other to allow interaction between them. As can be seen in the discussion below, each of the functional components of a replica state machine supports advantageous functionality based on one embodiment.<u style="single">Proposal management</u>
0039Each local application node 205 applies for a series of proposals to its respective replication state machine 200. The set of proposals submitted by each local node 6 constitutes a local sequence of each local node 205, which sequence can be maintained within the local sequencer 270 of each replication state machine 200. The proposal manager 220 of each replication state machine 200 subscribes each set of proposals into a global proposal sequence, one for each. This global proposed sequence can be maintained within the global sequencer 280 of each replication state machine 200. Global proposal sequences have the following properties: Each proposal in each local sequence occurs exactly once in each global sequence, and the relative order of any two proposals within the local sequence is: Optionally, the global sequence (with or without the local order saved) that is stored within each global sequence and associated with all local application nodes 205 is the same.
0040If a thread at local application node 205 applies for a proposal (eg, a write step) to each replication state machine 200, the replication state machine 200 assigns a local serial number to that proposal. The replication state machine 200 then determines the consent number for the proposal. As will become clear from the discussion below, the consent number determines the position of each proposal within the global sequence. The replication state machine 200 then stores a record of its proposals in its own fixed proposal storage 230. The replication state machine 200 then returns control of the threads of the local application node to the local application node. Therefore, the thread may be available for use by the local application rather than idle while the consent protocol is running. The replication state machine then invokes the consent protocol for its proposal via Consent Manager 240. At the end of the consent protocol, the replication state machine 200 compares the consent reached by the consent protocol with the proposed consent contained in the proposal. If the consent reached by the consent manager 240 is the same as that of the proposal, the replication state machine 200 terminates processing of the proposal. Otherwise, the replication state machine 200 repeatedly attempts to accept the proposal, using the new consent number, until the consent reached by the consent manager is the same as that of the proposal. Upon reaching the agreement conclusion, each local application node 205 queues what is now agreed to the proposal into its global sequence. Each local application node 205 of the distributed application then dequeues the proposals contained within the global sequence and executes them.
0041FIG. 4 shows an embodiment of the proposal based on one embodiment. This proposal is referred to herein as Proposal 300. Proposal 300 may include a proposer identifier 320 (ie, a local application node identifier), a local serial number (LSN) 330, a global serial number (GSN) 340, a consent number 350, and a proposal content 360. Preferably, but not necessarily, the proposals issued by each local application node have the structure of proposal 300.
0042FIG. 5 shows an embodiment of a local sequence based on one embodiment. This local sequence is referred to here as the local sequence 400. The local sequence 400 can include the content of each proposal for each local application node 205. More specifically, such content includes the proposer identifier, local serial number (LSN), global serial number (GSN), consent number, and proposal content. Preferably, but not necessarily, the local sequence associated with each replication state machine 200 has the structure of the local sequence 400.
0043FIG. 6 shows an embodiment of a global sequence based on one embodiment. This global sequence is referred to here as the global sequence 500. A global sequence can include a global sequence number for a set of suggestions and a local sequence handle. In one embodiment, the local sequence handle may be a pointer to each local sequence (ie, local sequence 400 as drawn). In other embodiments, the local sequence handle may be the key to the table of the local sequence. Preferably, but not necessarily, the global sequence associated with each replication state machine 200 has the structure of the global sequence 500. Concurrent consent
0044The replication state machine 200 drawn in FIGS. 2 and 3 is a replication state machine based on one embodiment, and a plurality of proposals from one proposer while optionally saving the order in which the proposals are submitted by the proposer. Includes a concurrent consent mechanism that allows consent to proceed concurrently. On the other hand, a conventional replication state machine attempts to agree to a proposal after the previous proposal has been agreed. This traditional replication state machine approach ensures that the traditional replication state machine preserves the proposed local order. Therefore, if the proposer proposes Proposal A first and then Proposal B, the conventional replication state machine agrees that Proposal A is agreed at the same time as Proposal B or before Proposal B. Guarantee. However, unlike the replication state machine that implements the backoff mechanism based on one embodiment, this conventional method does not activate the consent of proposal B until the consent of proposal A is reached, so that the consent of the conventional replication state machine is not activated. Slow down the operation.
0045With reference to aspects of one embodiment, each object (ie, entry) in the global sequence is numbered consecutively. The number associated with one object in a global sequence identifies its own position with respect to other objects in that global sequence. For example, an object numbered 5 precedes an object numbered 6 and precedes an object numbered 4. In addition, each object in the global sequence contains a handle to the local sequence, such as the local sequence handle 400 shown in Figure 5. If the application does not require storage in submission order (ie, source-issued order), each object in the global sequence contains the proposal itself. In this case, the suggestions can be obtained directly from the global sequence rather than indirectly via the local sequence. In one of several possible embodiments, the handle to the local sequence may be a pointer to the local sequence. In other embodiments, the handle to the local sequence may be the key to the table of the local sequence.
0046With reference to FIGS. 2 and 3, each local sequence contains a proposal for a replication state machine 200 proposed by one of the proponents of the replication state machine 200. Each local sequence node 205 of the replication state machine 200 maintains its own local sequence of the proposer associated with each replication state machine 200. Objects in the local sequence are continuously numbered. The number associated with an object in the local sequence identifies its own position with respect to other objects in that local sequence. For example, an object numbered 5 precedes an object numbered 6 and precedes an object numbered 4. Each object in the local sequence contains a proposal for replication state machine 200.
0047At each local application node 205 on the replication state machine 200, the proposal may be added to the global sequence after reaching consent to the proposal. Proposer identification (eg, proposer identifier 320 in FIG. 4) can be used as a key to find a local sequence from a table of local sequences. The proposal's local serial number (LSN) determines the position of the proposal within the local sequence. The proposal can then be inserted at a determined position in the local sequence. The proposal consent number (eg, consent number 350 in FIG. 4) determines the position of the proposal in the global sequence. Handles for local sequences can be inserted at determined positions in the global sequence (ie, based on consent numbers). The GSN is an optional bookkeeping field for associating a proposal with the purpose of specifying its actual position in the global sequence of the proposal when it is consumed, as described in the following paragraphs. ..
0048In one embodiment, a dedicated thread consumes the global sequence. The thread waits for the next position in the global sequence to be added. The thread then restores the local sequence stored at that location in the global sequence. The thread then waits for the next position in the local sequence to be added. The thread then restores the proposal for replication state machine 200 stored at that location in the local sequence. Those skilled in the art will appreciate that proposals do not necessarily have to be restored based on the order of consent numbers, but will be restored in exactly the same order on all application nodes. This restore order may be recorded in the GSN field for bookkeeping convenience, but is otherwise not essential to the operation of the replication state machine 200. For example, suppose application node (A) submits its first two proposals to a replication state machine (LSN1 and LSN2). Further assume that the replication state machine unexpectedly reaches LSN2 consent before reaching LSN1 consent. Therefore, assume that the consent number for A: 1 (LSN1 from application node A) is 27 and the consent number for LSN2 is 26 (ie, there is a total of 25 prior proposal consents from other application nodes. , A: 1 and A: 2 did not intervene with the consent of the proposals from other application nodes). Using the above method, A: 1 is restored from the global sequence at position 26 and A: 2 is restored from the global sequence at position 27. Thus, GSN adheres to the order of LSNs, but consent numbers do not necessarily have to. This technique allows a replica state machine based on one embodiment to process multiple consents at the same time.
0049The thread then applies the proposal for replication state machine 200. In one embodiment, the proposed application can be achieved by invoking a callback function registered by the application on the replication state machine 200.<u style="single">Backoff and collision avoidance</u>
0050A replication state machine based on one embodiment (eg, replication state machine 200) may include a backoff mechanism to avoid repeated proposer preempts in the consent protocol of consent manager 240. On the other hand, in a conventional replication state machine, when the round launched by the first proposer takes precedence over the round launched by the second proposer, the preferred proposer is higher than that of the predecessor. Allows the number to immediately launch a new round. Unfavorably, this conventional approach lays the groundwork for repeated round preemptions, causing the consent protocol to flutter for an unacceptably long time.
0051In the promotion of backoff based on one embodiment, the proposer calculates the duration of the preemption delay when the round is preempted. The proposer then follows the traditional algorithm for launching the next round and waits for the calculated duration before launching the next round.
0052In the promotion of collision avoidance based on one embodiment, when the first proposer senses that the round has been initiated by the second proposer, the first proposer determines the duration of the ongoing round delay period. calculate. The first proposer refrains from invoking the round until the calculated delay duration expires.
0053In one embodiment, the delay given increases exponentially with successive preemptions of the round. In addition, the delay is preferably determined randomly.
0054There are several methods that can be used to determine the duration of a given delay. One source of inspiration for viable methods is the literature on carrier sense multiple access / collision detection (CSMA / CD) for non-switched Ethernet. The CSMA / CD protocol is a set of rules that determine how a network device responds when two network devices attempt to use a data channel at the same time.
0055In some possible embodiments, the calculated delay duration is determined by the following method. The administrator who deploys the replication state machine 200 sets four numerical values. For the purposes of depiction of this embodiment, these values are referred to as A, U, R, and X. In an effective configuration, the value R is greater than 0 and less than 1, the value A is greater than 0, the value X is greater than 1, and the value U is greater than the value A. The execution time of the consent protocol can be estimated. One of the several possible estimates of the consent protocol execution time is a moving average of the consent protocol's past execution times. For the purposes of this discussion, this estimated value is referred to as E. The value M is the value obtained by multiplying A by U. The larger of the two values A and E is selected. For the purposes of this discussion, this selected value is referred to as F. The value C is the value obtained by multiplying F by X. The random value V is generated from a uniform distribution between 0 and C times R. If C is greater than M, then D is calculated by subtracting V from C. Otherwise, D is calculated by adding V to C.
0056The calculated value D can be used as an ongoing round delay. It can also be used as a preempt delay when the local application node 205 is preempted for the first time while an instance of the consent protocol is running. When the local application node 205 is preempted a second time or later during the execution of an instance of the consent protocol, the new value D is calculated using the old value D instead of the value A in the above method. The new value D can be used as a preemption delay.<u style="single">Reuse of fixed storage</u>
0057A replication state machine based on one embodiment (eg, replication state machine 200) reuses the fixed storage used to ensure fault tolerance and high utility. Referring to FIGS. 2 and 3, the storage reclaimer 290 determines the position of the replication state machine 200 in the global sequence of proposed proposals, and after all application nodes are notified of this position, the proposal store 230 Erase the record of the proposed proposal from. Each local application node 205 sends a message to each of the other local nodes 205 at regular intervals indicating the highest consecutive position added in the copy of its global sequence. Send. The storage reclaimer 290, at regular intervals, gives all consents up to the highest of consecutively added positions, no longer required by the local application node, from all copies of the global sequence. to erase. In this way, each replication state machine 200 reuses fixed storage.<u style="single">Weak booking</u>
0058A replication state machine based on one embodiment (eg, replication state machine 200) optionally provides a weak reservation mechanism to terminate the proposer's preemption under normal operating conditions. With reference to FIGS. 2 and 3, each proposer driving each replication state machine 200 is continuously numbered. For example, if there are three proposers, they are numbered 1,2,3. The proposer number determines which proposal of each replication state machine 200 will be driven by the corresponding proposer. If the number of the proposer is M and there are N proposers, the proposer is M + (k × N) (that is, M plus k times N, where k is 0 or more. It will drive multiple proposals numbered with). In order to allow such a system to evolve even when not all proposers of the distributed application system are responsive, if the proposals for the replication state machine are not determined in a timely manner, then each replication state. Any proposer associated with the machine 200 may propose "no operation" (ie no-op) for that proposal. To make this optimization transparent to the distributed application, the replication state machine 200 does not deliver no-op proposals for the distributed application. No operation is a computational step that generally has no effect and, in particular, does not change the state of the associated replication state machine.<u style="single">Distinguished and fair round numbers</u>
0059A replication state machine based on one embodiment ensures that one of a plurality of competing proposers is not preempted when using the same round number as the competing proposer. Traditional replication state machines, on the other hand, do not include a mechanism to ensure that one of a number of competing proposers is not preempted when using the same round number as the competing proposer. The round number in such a conventional replication state machine can be a monotonous value, which allows all proposers to be preempted.
0060In one embodiment, the round number may include a distinguishing component in addition to the monotonous component. In one embodiment, a small distinguishable integer may be associated with each proposer of each replication state machine 200. The distinguishable integers serve to resolve conflicts in the favor of the proposer with the highest distinguishing component. Round numbers include random components in addition to monotonous and distinguishing components. This type of round number ensures that one of multiple competing proposers is not preempted when using the same round number for competing proposals, and conflict resolution is specific among the proposers. Guarantee that one of the people will not be given preferential or cold treatment forever (ie, through the random component of the round number).
0061The mechanism for comparing two round numbers works as follows. Round numbers with larger monotonous components are said to be higher than others. If the monotonous components of the two round numbers are the same, then the round number with the larger random component is said to be higher than the other. If the random numbers are not distinguished by these two comparisons, then the round number with the larger distinguishing component is considered to be higher than the others. If the random numbers are not distinguished by these three comparisons, the two round numbers are considered equal.<u style="single">Efficient fixed storage reuse</u>
0062With reference to FIGS. 3 and 4, the record in the proposed fixed store 230 of the replication state machine 200 is composed of a plurality of groups. Each group stores a record of proposed proposals with a continuous local serial number 330. For example, the records with local serial numbers # 1 to # 10000 belong to group 1, the records with local serial numbers # 10001 to # 20000 belong to group 2, and so on.
0063With reference to the persistent proposal groups, each group can be stored in such a way that the storage resources used by the entire group are efficiently reused. For example, in a file-based storage system, each group uses its own file or set of files.
0064Looking at a group of more persistent suggestions, Storage Reclaimer 290 monitors individual record deletion requests, but does not delete individual records at the time of the request. When a cumulative delete request for an individual record contains all the records in a group, the storage reclaimer 290 efficiently reuses the storage resources used by that group. For example, in a file-based storage system, the files or filesets used by the group may be deleted.
0065The record in the consent store 245 of the replication state machine 200 is composed of a plurality of groups. Each group stores a record of instances of the consent protocol, along with a contiguous consent instance number 150. For example, the records of consent instance numbers # 1 to # 10000 belong to group 1, the records of consent instance numbers # 10001 to # 20000 belong to group 2, and so on.
0066By referencing a group of instances of the consent protocol, each group can be stored in such a way that the storage resources used by the entire group are efficiently reused. For example, in a file-based storage system, each group uses its own file or set of files.
0067Further referring to the group of instances of the consent protocol, the storage reclaimer 290 monitors individual record deletion requests, but does not delete individual records at the time of the request. When a cumulative delete request for an individual record contains all the records in a group, the storage reclaimer 290 efficiently reuses the storage resources used by that group. For example, in a file-based storage system, the files or filesets used by the group may be deleted.<u style="single">Efficient use of small suggestions</u>
0068Referring to FIGS. 3 and 4, a replication state machine based on one embodiment (eg, replication state machine 200) is a plurality of local application nodes 205 of multiple proposals applied to the replication state machine 200. Transmission from our proposal side to the receiving side of a plurality of local application nodes 205 is performed by batch processing. Such a practice is based on one embodiment in the context where the proposed size of the replication state machine is small compared to the size of the packet of data in the underlying packet-based communication protocol used by the replication state machine. Allows replication-state machines to efficiently utilize packet-based communication protocols.
0069In one embodiment, such batches of proposals may be treated as a single proposal by the consent protocol. According to this aspect, at each local node 205, each replication state machine 200 determines the consent number 350 of the first bundle of proposed proposals, while the proposals proposed at each local application node 205 Can accumulate in the second bundle of proposals. When the consent number 150 for the first bundle is determined, the replication state machine 200 invokes the decision for the consent instance number 350 for the second bundle, and the proposals proposed at the local application node 205 accumulate in the third bundle. The process is continued, and so on.<u style="single">Efficiently handle big proposal 110</u>
0070Can be performed in order to reduce network bandwidth for large proposals, proposals replication state machine in accordance with one embodiment is short id (e.g., 16 bytes of the glow can tag each suggested by the unique id Bal) And to encode the proposal into a format referenced as a file-based proposal, or one of these. The big proposal, on the other hand, causes the problem for traditional replication state machines that such large proposals are transmitted essentially multiple times when driven by the consent protocol of traditional replication state machines. .. Such multiple transmissions are not preferred, as the size of large proposals can be several megabytes or even gigabytes.
0071When sending a large proposal, in one embodiment, once the actual proposal has been successfully sent to the network endpoint, only the short proposal identifier is sent. File-based proposals essentially have a file pointer in memory, whereas the actual proposal content in the file is kept on disk. When carrying such file-based proposals over the network, the replication state machine based on one embodiment uses an efficient fault-tolerant file streaming protocol. Such transport is handled by DFTP layer 250 of the replication state machine 200 (Figure 3). DFTP Layer 250 monitors file-based proposals and network endpoints. DFTP Layer 250 ensures that the file-based protocol is sent only once to the endpoints of the network. In the event of a failure event leading to a partial transfer, the file-based proposal is read from any available endpoint that has the requested portion of the file.
0072In one embodiment, the DFTP implementation uses native transmit files or memory-mapped files for efficient file transfer, if the operating system supports these features. If the node requesting the file cannot reach the original caller, it looks for an alternative caller, another node in the system that happens to have the file. When operating over the TCP protocol, DFTP utilizes multiple TCP connections to take advantage of the optimal high-bandwidth connections that are also subject to high latency. In addition, the TCP protocol window size can be changed appropriately and / or desirable to take advantage of the high bandwidth connections that are also subject to high latency.
0073Let's move on to the discussion of extensible and active replication of information repositories. In one embodiment, the implementation of such replication based on one embodiment utilizes the replication state machine described above. More specifically, providing such a copy based on one embodiment has a positive effect on the scalability, reliability, and fault tolerance of such a copy state machine. Therefore, the execution of a replication state machine based on one embodiment has a positive effect on such replication in a distributed computer system architecture. In performing the replication of the Information Repository in one embodiment, all or part of the following objectives will be met:
0074With reference to FIG. 7, an embodiment of the Replicator 120 based on one embodiment is shown, which is referred to herein as the Replicator 600. The Replicator 600 consists of multiple functional modules including a Replicator Client Interface 610, a Prequalifier 620, a Replication State Machine 630, a Scheduler 640, a Replicator Repository Interface 650, a Result Handler 660, and an Administrator Console 670. Replicator Client Interface 610, Prequalifier 620, Replication State Machine 630, Scheduler 640, Replicator Repository Interface 650, Result Handler 660, and Administrator Console 670, respectively, to enable interaction between them, respectively. It is connected to each other with some of the other modules. The replication state machine 200 functionally discussed with reference to FIGS. 2 to 6 is an example of the replication state machine 630 of the replicator 600. Therefore, the replication state machine 630 is reliable, useful, extensible, and fault tolerant.
0075FIG. 8 shows an embodiment of the placement of the Replicator 600 in a multi-site computer texture based on one embodiment. A multi-site computer texture may include multiple distributed application systems 601. Each distributed application system 601 may include multiple clients 680, replicator 600, repository client interface 690, repository 695 (ie, information repository), and network 699. The network 699 generally does not have to be any one component of multiple distributed application systems 601 but is between the client 680 of each distributed application system 601 and its replicator 600, and It may be connected between the repository client interface 690 of each distributed application system 601 and its respective replicator 600, so that the interaction of the client 680, replicator 600, and repository 695 of each distributed application system 601 It will be possible to connect these to each other to enable. The network may also be connected between replicators 600 of all distributed application systems 601 so that interactions are possible between replicators 600 of all distributed application systems 601. .. Multiple networks 699 may be isolated from each other, but they do not have to be. For example, the same network can fulfill all three roles disclosed above.
0076As shown in FIG. 8, there are three clients 680 "near" each of the multiple repositories 695 (ie, the system elements of the multiple distributed application systems 601 each containing the repository 695). By being close, a particular one in multiple repositories 690 A particular one in multiple client 680s in the vicinity chooses to access that particular one in multiple repositories 690. It means deaf. Alternatively, that particular one of the multiple clients 680 could potentially access its repository 695 in any distributed application system 601.
0077Operators of a distributed computer system based on an embodiment include users of client 680 and one or more administrators of distributed application system 601. The client 680 user follows the instructions in the client user manual. Since many advantageous aspects of an embodiment can be transparent to the user, the user may be unaware of the fact that he is using a replicator based on one embodiment. The administrator will configure the network appropriately as needed and, if necessary for the operation, in addition to the standard tasks of managing the Repository 695 itself.
0078The replication state machines 630 of each distributed application system 601 communicate with each other over the network 699. Each Replicator repository interface 650 interacts with the repository 695 of its respective distributed application system 601 through network 699. Client 680 interacts with Replicator Client Interface 610 through Network 699. Optionally, if a particular distributed application system 601 containing a particular client 680 is not available at any given time to provide the requested functionality, then that particular client of that particular distributed application system 601. It is also possible to use a product such as Cisco Systems Director to allow the 680 to fail over to any of the other distributed application systems 601.
0079With reference to FIGS. 7 and 8, the replicator client interface 610 controls an interface with a particular one of a plurality of client 680s (ie, a particular client 680) associated with the target repository 695. Replicator client interface 610 reconfigures a command issued on network 699 by a particular client 680 and delivers that command to prequalifier 620. The Prequalifier 620 enables efficient operation of the Replicator 600, but is not essential for the useful and advantageous operation of the Replicator 600.
0080For each command, the prequalifier 620 optionally determines if the command is destined to fail, and if so, an appropriate error message or error condition to the particular client 680 above. Decide to be returned. If so, the error message or error status can be sent back to the replicator client interface 610, which delivers the error message or error status to the particular client 680 described above. Since then, the command may not be processed by any other Replicator 600.
0081For each command, the prequalifier 620 optionally determines whether the command can bypass the replication state machine 630, or both the replication state machine 630 and the scheduler 640. If the prequalifier 620 does not determine that the replication state machine 630 can be bypassed, the command may be delivered to the replication state machine 630. The replication state machine 630 collates all commands submitted to itself and its fellow replication state machines 630 in each of the other associated Replicator 600s of the distributed application system 601. The sequence of operations may be guaranteed to be the same in all distributed application systems 601. In each of the plurality of distributed application systems 601, each replication state machine 630 sequentially delivers the commands collated as described above to the respective scheduler 640.
0082The scheduler 640 performs dependency analysis on the commands delivered to it and determines the weakest partial ordering, which still guarantees one-copy serializability. Such dependency analysis and one-copy serialization can be found in Wesley Addison's prior art literature entitled "Concurrent Control & Recovery in Database System". It is disclosed and published in a reference book by P. Berstein. The scheduler 640 is then allowed by the constructed partial order, while otherwise continuously delivering commands to the Replicator repository interface 650.
0083Replicator repository interface 650 delivers commands to repository 695. Correspondingly, one of three consequences occurs. The replicator repository interface 650 then delivers the results that are occurring to the result handler 660.
0084The first of the above three results may include the repository 695 returning a response to the command. This response includes results, states, or both, indicating that nothing went wrong during the execution of the command. If the origin of the command is local, the result handler 660 delivers the response to the replicator client interface 610, and the replicator client interface 610 delivers the response to the client 680. If the command originates from another distributed application system 601 replicator, the response is preferably dropped.
0085The second of the above three results may include repository 695 returning an error condition. The result handler 660 determines whether the error condition is a deterministic error in the repository 695 (ie, whether the same or equivalent error occurs in each of the other distributed application systems 601). decide. If the error determination is ambiguous, the result handler 660 attempts to compare the error with the result in another distributed application system 601. If this does not disambiguate, or if the error is not clearly deterministic, the result handler 660 aborts the operation of Replicator 600 and through the admin console 670 (ie, administer). Notify the operator (through issuing notifications via the console 670).
0086If the replicator is a CVS replicator, the result handler uses a list of error patterns to flag deterministic errors, as discussed below for CVS-specific functionality. The result handler 660 uses these patterns to perform regular expression matching in the response stream.
0087The third of the above three results may include a hung (ie, never return from command execution) repository 695. In one embodiment, this result can be treated just like a non-deterministic error, as discussed with reference to the second of the three results above.
0088Based on one embodiment, each Replicator 600 may be configured alternative. In an alternative embodiment, the Replicator 600 can be embedded in the Client 680 of Repository 695 and driven directly. In other alternative embodiments, the Replicator 600 may be embedded within the client interface 690 to Repository 695. In other alternative embodiments, the Replicator 600 may be embedded within Repository 695. In other alternative embodiments, the replicator's global sequencer may be based on other techniques, with a compromise of associated robustness and quality of service. One of the few possible examples of such techniques is group communication. In another alternative embodiment, the Replicator 600 drives one or more Repositories 695, with the compromising robustness and quality of service that accompanies it. In other alternative embodiments, the Replicator 600 modules are integrated into coarser modules, split into finer modules, or both. In other alternative embodiments, the information contained in repository 695 of each distributed application system 601 is applied to each of the other distributed application systems 601 as a redundant precaution against deviations from one-copy serialability. The responses of all distributed application systems 601 are compared to ensure that they remain consistent.
0089Referring to FIGS. 7 and 8, each of the plurality of repositories 695 discussed above may be a Concurrent Versions System (CVS) repository and the client 680 may be a corresponding CVS client. If Repository 695 is a CVS repository and Client 680 is a CVS Client, the interfaces associated with Repository 695 and Client 680 are CVS-specific interfaces (eg, Replicator CVS Interface, Replicator CVS Repository Interface, and Repository CVS Client). Interface). Further, based on one embodiment, the Replicator 600 can be modified to have features specifically and specifically configured for use with the CVS repository.
0090The replicator client interface 610 disclosed herein may be specifically configured to interface a CVS client with a target CVS repository. To do this, the Replicator Client Interface 610 stores the bytes arriving from the CVS client in a memory-mapped file buffer. Replicator client interface 610 detects the end of a CVS command when it finds a valid command string in the arriving byte stream. An unrestricted list of such valid command strings is, but is not limited to, "Root," Valid-responses "," valid-requests "," Max-dotdot "," Static-directory "," Sticky, Entry, Kopt, Checkin-time, Modified, Is-modified, UseUnchanged, Unchanged, Notify, Questionable, Argument, Argumentx , "Global_option", "Gzip-stream", "wrapper-sendme-rcsOptions", "Set", "expand-modules", "ci", "co", "update", "diff", "log", " rlog, "list", "rlist", "global-list-quiet", "ls", "add", "remove", "update-patches", "gzip-file-contents", "status", " rdiff, "tag", "rtag", "import", "admin", "export", "history", "release", "watch-on", "watch-off", "watch-add", " watch-remove, "watchers", "editors", "init", "ann"
0091The Replicator Client Interface 610 then attempts to classify the arriving CVS commands into read or write commands. An unlimited list of valid write command strings, but not limited to, "ci", "tag", "rtag", "admin", "import", "add", "remove", "watch" Can include "-on", "watch-off", and "init". All commands in a valid command string that are not in the list of valid write command strings are considered read here for the list of valid command strings.
0092Read commands are delivered directly to the CVS Replicator repository interface for execution by the target CVS repository. The CVS write command is optionally delivered to the prequalifier module 20.
0093For each CVS write command, Prequalifier 20 can optionally determine if CVS is destined to fail, and if so, the appropriate error to be returned to the CVS client. Determine the message or error condition. Failure detection can be based on matching the results or state byte streams returned by the CVS repository with known error patterns. Examples of known system error patterns are, but are not limited to, cannot generate symbolic links from. * To. *, Cannot start the server via rsh, cannot fstat. *, Generate temporary files. Failed, cannot open dbm file. * For generation, cannot write to. *, Cannot stat history file, cannot open history file:. *, Cannot open ". *", For mapping Unable to stat RCS archive. *, Unable to open file. * For comparison, exhausted virtual memory, unable to ftello in RCS file, unreadable. *, List of auxiliary groups Can't get, can't fsync file. * After copying, can't stat. *, Can't open current directory, can't stat directory. *, Can't write to. *, Can't readlink. *, Can't readpipe Could not close, could not change to directory. *, Could not generate temporary files, could not get file information for. *, Could not open diff output file. *. Can't generate. *, Can't get working directory, Can't lstat. *, Fork for diff failed for. *, Couldn't get information for. *,. * Cannot change mode due to. * Cannot ftero due to. *, Message verification failed, Temporary files cannot be stats, Out of memory, Directory ". *" in directory ". *" * Cannot be opened,. * Could not be stat, directory. * Cannot be opened, fwrite failed, temporary file. * Cannot be generated, temporary file cannot be stat,. * Cannot be stat, ". *" Unable to read, error when diffing. *, Unable to generate special file. *, Unable to close history file:. *, Unable to map memory to RCS archive *, generate directory ". *" Can't read file. * For copy, can't generate pipe, can't open temporary file. *, Can't move file. *, Can't open, can't seek end of history file, can't chdir to. * Included: Length read failed,. * Cannot be executed,. * Cannot be fdopened, and temporary file size cannot be found. Examples of known non-system error patterns are, but not limited to, internal errors: no such repository, no desired version found, getsockname failed, warning: rewriting RCS file Ferror set in between, internal error: islink does not like readlink, access is denied, device files of this system cannot be compared, server internal error: unprocessed case in server update, receive signal. Internal error: No revision information for it, Protocol error: Duplicate mode, Server internal error: No mode in server update, rcsbuf cache opens: Internal error, Fatal error, Abort, Fatal error: Exit ,. *: Unexpected EOF,. *: Confused revision number, invalid rcs file, EOF on key in RCS file, RCS file in CVS always ends with v, hard link information lost,. * : Unable to read end of file, open rcsbuf: Internal error, out of memory, unable to allocate infopath, unexpected dying gasp from. *, Internal error: incorrect date. *, Kerberos authentication failed:. *,. * Delta. *: Unexpected EOF, RCS file Unexpected EOF when reading. *, ERROR: Out of space-suspended, flow control EOF, unable to fseeko RCS file. *, Checksum failure for. *, CVS internal error: unknown state \ d +, internal error: Illegal argument for run_print, unable to copy device files to this system, unexpected file termination when reading. *, Out of memory, internal error: no RCS file parsed, internal error: premature EOF in RCS_copydeltas , Internal error: Testing support for unknown response \ ?. EOF in values in RCS file. *, PANIC \ * management file is lost \ !, end of file too early to read. *, EOF while looking for values in RCS file. * , Unable to continue, Read lock failed-Give up, Unexpected EOF when reading. *, Unable to resurrect ". *", RCS file deleted by second party, Your apparent username. * Unknown to this system, file attribute database corruption, loss of tabs in. *, device files cannot be imported into this system,. *: special files of unknown type cannot be imported, *: unknown type Unable to import special files, error: mkdir. *-Not added, unable to generate write locks in repository ". *", Unable to generate. *, Unable to generate special files on this system ,. * Cannot save, device file cannot be saved to this system, error when parsing repository file. *, File may be corrupted, file.
0094As discussed with reference to FIGS. 7 and 8, for each command, the prequalifier 620 is destined to fail and can bypass both the replication state machine 630 and the scheduler 640. Can be determined. In the case of CVS-specific functionality, if the prequalifier 620 does not determine that the replication state machine 630 can be bypassed, the command can be converted to a CVS proposed command. In addition to the actual CVS command byte array, the CVS Proposed Command contains a lock set that represents a write lock that would be acquired by the CVS repository if the CVS Proposed Command were executed directly by the CVS repository. The scheduler 640 makes use of this lock set, as discussed below.
0095The CVS proposal command can be delivered to the replication state machine 630. The replication state machine 630 collates all commands submitted to itself and its companions, the replication state machine 630, in each of the other replicators into a single sequence. This sequence is guaranteed to be the same for all replicas. In each of the distributed application systems 601, the replication state machine 630 sequentially delivers the commands collated as described above to the scheduler 640.
0096The scheduler 640 performs dependency analysis on the commands delivered to it and determines the weakest partial ordering, which still guarantees one-copy serializability. The scheduler 640 delivers commands to the CVS Replicator repository interface, otherwise continuously, as permitted by the constructed partial order.
0097Based on one embodiment, dependency analysis can be based on testing lock collisions. Each CVS proposal command submitted to the scheduler contains a lock set. The scheduler ensures that a command is delivered to the CVS repository only if its lockset does not conflict with the lockset of another command. If a collision is detected, the command waits on the queue and is scheduled at a later point in time than when all locks in the lock set were obtained without collision.
0098As mentioned above, the implementation of the multi-site computer system architecture has a positive impact on the scalability, reliability, usefulness, and fort tolerance of such replica state machines. Efficient scaling requires an efficient process to add newly allocated application nodes (or simply nodes) to the system. However, in order for newly added nodes to be able to participate in a distributed computer system, they must be given a certain amount of information. For example, a new node must be qualified to participate in a collaborative project, and an existing location and a newly invited node that can be seen from it interacts with the exchange of messages. Must be talked about with the allowed nodes. According to one embodiment, they correspond to the messaging model, the node induction method, which is effective to allow the inducer node to bring the newcomer node to a decentralized computer system. It is made possible by the equipment and systems and by allowing the deployed nodes to do useful work.<u style="single">Messaging model</u>
0099Here, the term "leader" or "leader node" means a node that at least initiates the introduction of another node, i.e., a "newcomer node" into a distributed computer system. According to one embodiment, the inducer node and the newcomer node communicate with each other by transmitting a message using an asynchronous and uncomplicated model as described below. -One of the processes runs at any speed and may fail and restart by stopping, -Some information must be remembered (ie, persistent) before and after the reboot, as the process may fail at any point. -Messages can take any length of time to be delivered and can be duplicated and lost, but not corrupted (corrupted messages are treated the same as undelivered messages and are therefore discarded by the recipient ).
0100FIG. 9 illustrates one aspect of a device, method, and system that enables the secure and authenticated introduction of a node into a group of nodes according to one embodiment. As shown therein, and according to one embodiment, such as the pre-authentication phase, the startup phase of the newcomer node, the placement phase of the bootstrap membership, and the recognition of the newcomer node and location. The method may include a method of introducing a node into a distributed computer system, and the system and equipment may be configured to perform. In addition, multiple post-deployment tasks can be performed. Each of these phases is described in detail below. A. <u style="single">Pre-authentication phase</u>
0101According to one embodiment, the pre-authentication phase can be performed before the newcomer node 206 is started, allowing the administrator 202 to be used in the deployment process and to preconfigure the deployment process. It can provide an opportunity to generate an implementation task that contains information to be done. Therefore, according to one embodiment, the pre-authentication phase can proceed without human interference. A.0 <u style="single">Generate new installation task</u>
0102Before the newcomer node 206 is launched, a deployment task may be generated on the inducer node 204 that contains the information required for a successful and complete deployment process. The use of persistent tasks is such that the information required within the deployment process and the state of the deployment process remain after the reboot of the inducer node 204, and the same deployment task is copied in another deployment ( Allow this information to be stored in the same location for reuse (cloned).
0103According to one embodiment, the introductory task may be configured to include three elements: an introductory ticket, a set of nodes that the newcomer node 206 should know about, and a set of post-introduction tasks. It is understood that other elements may be added and that other elements may replace these three elements. A.1 <u style="single">Introductory ticket</u>
0104The deployment is generated, for example, by an administrator and sent to newcomer node 206, as shown in B21 of FIG. This deployment ticket, for example, packages the detailed contacts of the inducer node 204 to the administrator 202, identifies (and controls) the details of the new node, and also some other platform or Provides a mechanism for identifying application configuration parameters for the new node. According to one embodiment, the introductory ticket may include: -Identification of deployment tasks -Node and location identification of newcomer node 206 -Location identification, host name, and port of inducer node 204 (basic information required for newcomer node 206 to access inducer node 204), and / or -Other, optional, platform / application configuration information
0105The introductory ticket is the same or functional as the newcomer node 206 will be able to see, see, and communicate with other selected nodes in the distributed computer system. May contain other information that achieves similar results. The introductory ticket can be configured, for example, as a file. To increase security, such deployment tickets can therefore be co-designed with the private key of the inducer node in the PKI system. Similarly, the newcomer node 206 is configured to use the public key of the inducer node 204 to verify the authenticity of the details contained in the introductory ticket. The implementations described and presented here are not limited to the PKI model of security, so other authentication and authority-defined methods can be used to their advantage. According to one embodiment, the introductory ticket is then sent out of band to engineer 208 performing the installation of the ink ducty node 206. According to one embodiment, the introductory ticket remains with the inducer node 204 and is "pushed" to the newcomer node 206 when the newcomer node 206 is activated. A.2 <u style="single">A set of nodes that newcomer nodes should be aware of</u>
0106According to one embodiment, the deployment task may include details of existing nodes to which the newcomer node 206 should be notified during the deployment process. It should be noted that the newcomer node 206 can communicate / work with the newcomer node 206 and / or be informed of other nodes in the permissible distributed computer system. is there. Such information is therefore favorably identified before the implementation process is initiated, if there is no human intervention. The selection of node or multiple nodes that newcomer node 206 is able to communicate with or is allowed to do is within the complete set or subset of existing nodes already deployed in the network of multiple nodes. It can be done with a user interface (UI) that allows the administrator 202 to select a set of nodes. This information is stored within the deployment task for later access. The UI may include, for example, a browser or mobile device application. A.3 <u style="single">Post-introduction tasks</u>
0107According to one embodiment, the post-introduction task may include a plurality of task details. One or more of such tasks may apply to a new newcomer node 206 following the implementation, for example joining an existing membership. It should be noted that this set of tasks may be empty if newcomer node 206 is not required to do anything when following the deployment. Once the introductory task is created and made persistent (eg, stored in non-volatile memory), the newcomer node 206 can be started. B. <u style="single">Launch newcomers</u>
0108According to one embodiment, the newcomer node 206 may be activated. B.1. <u style="single">If there is no introduction ticket on the newcomer node</u>
0109According to one embodiment, if the introductory ticket does not exist on the newcomer node 206, the newcomer node 206 is started or started with the basic configuration as described below. Wait for the inducer node 204 to come in contact with the bootstrap membership details that the newcomer node 206 will be a member of (ie, listen as shown in B22). B.2. <u style="single">If there is an introduction ticket on the newcomer node</u>
0110According to one embodiment, if the introductory ticket really exists on the newcomer node 206 at the time of activation shown in B23, the newcomer node 206 can be configured as follows: a) Analyze the information in the deployment task (optionally authenticate if necessary). b) Use this information to configure the application platform. c) Use this information to generate and turn on the Bootstrap Membership Request beacon shown in B24. The bootstrap membership request beacon may be configured to notify the inducer node 204 that the newcomer node 206 is initiating the deployment process, as shown in B25. According to one embodiment, a beacon is a process that is configured to repeatedly broadcast a message to a predetermined list of target recipients and removes the target recipients from the predetermined list. The message is broadcast to a given list until response acknowledgement is received from the individual target recipients. According to one embodiment, the bootstrap membership request may be configured to include deployment task identification, newcomer node and location identification, hostname, and port.
0111According to one embodiment, in response to the inducer node 204 receiving a bootstrap membership request from the newcomer node 206, the inducer node 204 relatives to the newcomer node 206 as shown in B26. A Bootstrap Membership Response may be returned, thereby disabling the request beacon as shown in B27. The inducer node 204 can then examine the deployment task to see if its node and location identification matches what was identified in the previously received deployment ticket, as shown in B28. it can. If this verification fails, i.e. the identification of its node and / or location does not match that in the deployment task, the inducer node 204 issues a Bootstrap Membership Denied message, as shown in B29. Can beacon to newcomer node 206.
0112Upon receiving the bootstrap membership denial message, the newcomer node 206 may be configured to send a BootstrapMembershipAck message in response and terminate processing as shown in B30. When the inducer node 204 receives the bootstrap membership ack message from the newcomer node 206 as shown in B31, it may invalidate the bootstrap membership denial message as shown in B32. C. <u style="single">Placement of bootstrap membership</u>
0113According to one embodiment, if the newcomer node 206 boots without a deployment ticket and the administrator 202 initiates the deployment process on the inducer node 204, or if the deployment task lookup is successful, bootstrap. Membership generation and placement can be performed using the following process:
0114According to one embodiment, the inducer node 204 performs the following according to one embodiment: 1. Generate bootstrap membership with the following a ~ c. Deterministically generated membership identification b. Inducer node 204 in the role of consent proposer and consent acceptor c. Newcomer node in learner role 2. Place membership as shown in B33. 3. Generate a deterministic state machine with reference to bootstrap membership, as shown in B34. 4. Beacon a Bootstrap Membership Ready message to newcomer node 206, as shown in B35.
0115According to one embodiment, upon receiving the bootstrap membership ready message as shown in B36, the newcomer node 206 performs the following according to one embodiment: 1. Generate bootstrap membership with the following a ~ c. a) Deterministically generated membership identification b) Inducer node 204 in the role of consent proposer and consent acceptor c) Newcomer node in learner role 2. Place membership as shown in B37. 3. Generate a deterministic state machine with reference to bootstrap membership, as shown in B38. 4. Beacon a bootstrap membership ack message to newcomer node 206, as shown in B39. According to one embodiment, when the inducer node 204 receives the bootstrap membership ack message, it should disable the bootstrap membership ready beacon, as shown in B40.<u style="single">D. Newcomer node and location awareness</u>
0116Following the placement of bootstrap membership, the newcomer node 206 is informed of the nodes and locations it should be aware of. This is achieved using the following process, according to one embodiment. 1. The inducer node 204 examines the deployment task to determine which location and node the newcomer node 206 should be informed of. 2. The deployment task returns a list of locations and nodes for this newcomer node 206. 3. The inducer node 204 proposes a set of nodes and locations to the deterministic state machine. 4. If consent is formed as shown in B41, the newcomer node 206 learns about the locations and nodes that need to be known, as shown in B42. 5. Following the newcomer node 206 learning about the node and location, the deployment process is complete.<u style="single">E. Post-introduction tasks</u>
0117Following node and location consent, or completion of the deployment process, it is possible to perform task installations identified within the deployment task. These tasks include, for example, creating a new membership that includes a newly introduced node, joining an existing membership (ie, including a newly introduced node within an existing membership). , Enforce membership changes), can include deploying and synchronizing duplicated entities.
0118FIG. 10 shows a block diagram of a computer system 1000 in which an embodiment can be implemented. The computer system 1000 may include a bus 1001 or other communication mechanism for communicating information and one or more processors 1002 that combine with the bus 1001 for processing information. Computer system 1000 further combines random access memory (RAM) or other dynamic storage device 1004 (referred to as main memory) with bus 1001 to store information and instructions executed by processor 1002. May include. Main memory 1004 can also be used to store temporary variables and other intermediate information during the execution of instructions by processor 1002. Computer system 1000 may also include read-only memory (ROM) and / or other static storage device 1006 that couples with bus 1001 to store information and instructions for processor 1002. A data storage device 1007, such as a porcelain disk or flash memory, may be coupled to bus 1001 to store information and instructions. The computer system 1000 may also be coupled to the display device 1010 via bus 1001 to display information to the user of the computer. An alphanumeric input device 1022, including alphanumeric characters and other keys, may be coupled to bus 1001 for information transmission and instruction selection to processor 1002. Other types of user input devices include cursor controls for transmitting instructional information to processor 1002 and selecting instructions, such as a mouse, trackball, and cursor instruction keys, and for controlling cursor movement on display 1021. It is 1023. The computer system 1000 may be coupled to the network 1026 and one or more nodes of the distributed computer system via a communication device (eg, a modem or NIC).
0119Embodiments relate to the use of computer systems and / or multiple such computer systems that introduce nodes into distributed computer systems. According to one embodiment, the methods and systems described herein may be provided by one or more computer systems 1000, depending on processor 1002, which executes a series of instructions contained within memory 1004. Such instructions can be read into memory 1004 from another computer-readable medium, such as data storage device 1007. Execution of a series of instructions contained in memory 1004 causes processor 1002 to perform the steps described herein and have functionality. In alternative embodiments, hard-wired circuits may be used in place of or in combination with software instructions to perform the embodiments. Therefore, embodiments are not limited to a particular combination of hardware circuits and software. In fact, it should be understood by those skilled in the art that any suitable computer system may perform the functions described herein. A computer system may include one or more microprocessors operating to perform the desired function. In one embodiment, instructions executed by one microprocessor or multiple microprocessors serve to cause the microprocessor to perform the steps described herein. The instructions may be stored on any computer readable medium. In one embodiment, it may be stored in a non-volatile semiconductor memory outside the microprocessor or embedded in the microprocessor. In other embodiments, the instructions are stored on disk and can be read into volatile semiconductor memory before being executed by the microprocessor.
0120Although some embodiments of the disclosure have been described, these embodiments are expressed as examples only and are not intended to limit the scope of the disclosure. In fact, the novel methods, devices, and systems described herein can be embodied in a variety of other forms. Moreover, in the forms of methods and devices described herein, various omissions, substitutions and modifications can be made without leaving the spirit of disclosure. The appended claims and their equivalents are intended to cover such forms or modifications as included in the scope and spirit of the disclosure. For example, one of ordinary skill in the art will appreciate that in various embodiments, the actual physical and logical structures may differ from those shown in the drawings. Depending on the embodiment, some steps described in the above example may be removed and other steps may be added. Also, the characteristics and attributes of the particular embodiments disclosed above can be mixed in different ways to form additional embodiments. All additional embodiments referred to herein are within the scope of the present disclosure. Although this disclosure provides some embodiments and applications, other embodiments, including embodiments that are apparent to those skilled in the art and do not provide all the features and benefits described herein, are also described in this publication. It is within the scope of disclosure. Therefore, the scope of this disclosure is intended to be defined only by reference to the appended claims.
10 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6 Sheet 7 Sheet 8 Sheet 9 Sheet 10
Every citation, both ways
| Document | Relation | Office |
|---|---|---|
| JP2007528557A | Cites | Japan |
| US20040254977A1 | Cites | United States of America |
97 members in 12 offices
Members97
| Document | Office | Kind | |
|---|---|---|---|
| US2006155729A1 | United States of America | A1 | |
| WO2006076530A2 | World Intellectual Property Organization (WIPO) | A2 | |
| WO2006076530A3 | World Intellectual Property Organization (WIPO) | A3 | |
| US2008140726A1 | United States of America | A1 | |
| US8364633B2 | United States of America | B2 | |
| CA2893302A1 | Canada | A1 | |
| CA2893327A1 | Canada | A1 | |
| US2014188971A1 | United States of America | A1 | |
| US2014189004A1 | United States of America | A1 | |
| WO2014105247A1 | World Intellectual Property Organization (WIPO) | A1 | |
| WO2014105248A1 | World Intellectual Property Organization (WIPO) | A1 | |
| CA2922665A1 | Canada | A1 | |
| US2015067002A1 | United States of America | A1 | |
| US2015067004A1 | United States of America | A1 | |
| WO2015031755A1 | World Intellectual Property Organization (WIPO) | A1 | |
| AU2013368486A1 | Australia | A1 | |
| AU2013368487A1 | Australia | A1 | |
| US2015278244A1 | United States of America | A1 | |
| CA2938768A1 | Canada | A1 | |
| WO2015153045A1 | World Intellectual Property Organization (WIPO) | A1 | |
| EP2939446A1 | European Patent Office (EPO) | A1 | |
| EP2939447A1 | European Patent Office (EPO) | A1 | |
| US2016019236A1 | United States of America | A1 | |
| US9264516B2 | United States of America | B2 | |
| JP2016505977A | Japan | A | |
| AU2014312103A1 | Australia | A1 | |
| JP2016510448A | Japan | A | |
| US9332069B2 | United States of America | B2 | |
| US9361311B2 | United States of America | B2 | |
| US2016191622A1 | United States of America | A1 | |
| EP3039549A1 | European Patent Office (EPO) | A1 | |
| AU2015241457A1 | Australia | A1 | |
| EP2939446A4 | European Patent Office (EPO) | A4 | |
| EP2939447A4 | European Patent Office (EPO) | A4 | |
| US9424272B2 | United States of America | B2 | |
| JP2016530636A | Japan | A | |
| US9467510B2 | United States of America | B2 | |
| US9495381B2 | United States of America | B2 | |
| US2017024411A1 | United States of America | A1 | |
| US2017026465A1 | United States of America | A1 | |
| GB201621753D0 | United Kingdom | D0 | |
| EP3127018A1 | European Patent Office (EPO) | A1 | |
| AU2013368487B2 | Australia | B2 | |
| AU2013368486B2 | Australia | B2 | |
| EP3039549A4 | European Patent Office (EPO) | A4 | |
| US2017193002A1 | United States of America | A1 | |
| JP2017519258A | Japan | A | |
| EP3127018A4 | European Patent Office (EPO) | A4 | |
| US9747301B2 | United States of America | B2 | |
| US9846704B2 | United States of America | B2 | |
| GB201719095D0 | United Kingdom | D0 | |
| US9900381B2 | United States of America | B2 | |
| JP6282284B2 | Japan | B2 | |
| JP6342425B2This record | Japan | B2 | |
| CA3077113A1 | Canada | A1 | |
| US2018181730A1 | United States of America | A1 | |
| WO2018114976A1 | World Intellectual Property Organization (WIPO) | A1 | |
| GB2557970A | United Kingdom | A | |
| JP6364083B2 | Japan | B2 | |
| EP3039549B1 | European Patent Office (EPO) | B1 | |
| ES2703901T3 | Spain | T3 | |
| US10268808B2 | United States of America | B2 | |
| US2019243954A1 | United States of America | A1 | |
| AU2014312103B2 | Australia | B2 | |
| AU2015241457B2 | Australia | B2 | |
| AU2019236685A1 | Australia | A1 | |
| EP3559844A1 | European Patent Office (EPO) | A1 | |
| US10481956B2 | United States of America | B2 | |
| MX2019007393A | Mexico | A | |
| BR112019012953A2 | Brazil | A2 | |
| EP2939447B1 | European Patent Office (EPO) | B1 | |
| KR20190139831A | Republic of Korea | A | |
| CN110603537A | China | A | |
| CA2893302C | Canada | C | |
| JP6628730B2 | Japan | B2 | |
| EP2939446B1 | European Patent Office (EPO) | B1 | |
| CA2938768C | Canada | C | |
| ES2768325T3 | Spain | T3 | |
| ES2774691T3 | Spain | T3 | |
| JP2020522083A | Japan | A | |
| US10783224B2 | United States of America | B2 | |
| US10795863B2 | United States of America | B2 | |
| GB2557970B | United Kingdom | B | |
| CA2893327C | Canada | C | |
| AU2019236685B2 | Australia | B2 | |
| US2021042266A1 | United States of America | A1 | |
| EP3127018B1 | European Patent Office (EPO) | B1 | |
| CA2922665C | Canada | C | |
| US2021326415A1 | United States of America | A1 | |
| ES2881606T3 | Spain | T3 | |
| US2022121623A1 | United States of America | A1 | |
| JP7265987B2 | Japan | B2 | |
| KR102533544B1 | Republic of Korea | B1 | |
| EP3559844B1 | European Patent Office (EPO) | B1 | |
| CN110603537B | China | B | |
| US11853263B2 | United States of America | B2 | |
| ES2963168T3 | Spain | T3 |
17 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Written notification of registration of transferJAPANESE INTERMEDIATE CODE: R350R350 | R350 | |
| Written request for registration of change of nameJAPANESE INTERMEDIATE CODE: R313533S533 | S533 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Receipt of annual feesJAPANESE INTERMEDIATE CODE: R250R250 | R250 | |
| Certificate of patent or registration of utility modelJAPANESE INTERMEDIATE CODE: R150R150 | R150 | |
| First payment of annual fees (during grant procedure)JAPANESE INTERMEDIATE CODE: A61A61 | A61 | |
| Written decision to grant a patent or to grant a registration (utility model)JAPANESE INTERMEDIATE CODE: A01A01 | A01 | |
| Decision of grant or rejection writtenTRDD | TRDD | |
| Request for written amendment filedJAPANESE INTERMEDIATE CODE: A523A521 | A521 | |
| Written request for extension of timeJAPANESE INTERMEDIATE CODE: A601A601 | A601 | |
| Notification of reasons for refusalJAPANESE INTERMEDIATE CODE: A131A131 | A131 | |
| Report on retrievalJAPANESE INTERMEDIATE CODE: A971007A977 | A977 | |
| Written request for application examinationJAPANESE INTERMEDIATE CODE: A621A621 | A621 |
Numbers
- Publication
- 6342425
- Application
- 2015550381
Titles2
- Japanese
- 分散型コンピューター環境におけるノードのグループへのノードの安全かつ認証された導入を可能にする方法、装置、及びシステム
- English
- Methods, devices, and systems that enable secure and authenticated deployment of nodes into groups of nodes in a distributed computer environment.
Classification
- CPC, 12
- H04L67/1095
- H04L67/34
- H04L69/24
- H04L12/1886
- G06F11/2094
- G06F11/2097
- G06F2201/855
- G06Q10/10
- G06Q10/101
- H04L67/10
- H04L67/133
- G06F9/5083
- IPC, 3
- G06F8 60
- G06F9 4401
- G06F9 50
